Skip to main content

· 11 min read
百岁

本文面向刚加入 TIS 团队的开发工程师。读完本文,你将理解"伴生操作类"是什么、为什么需要它、以及如何从零开发一个完整的 Manipulate 插件。


一、背景:什么是"伴生操作类"?

TIS 的核心实体(如数据管道 DataXProcessor、本体域 OntologyDomain)本身只负责描述配置,不直接处理用户在界面上触发的"操作"——比如克隆、开启某个功能、触发推断任务等。

这些"操作"由伴生操作类(Manipulate)承担。每个核心实体可以挂载若干个 Manipulate 插件,每个插件对应一种用户可触发的操作。

核心实体(如 OntologyDomain)
└── 伴生操作类 1:EnableChatBI(开启智能问数)
└── 伴生操作类 2:InferOntologyFromLLM(LLM 推断本体关系)
└── 伴生操作类 N:...

用户在界面上点击某个操作按钮时,TIS 框架会找到对应的 Manipulate 实例,调用其 manipuldateProcess 方法。


二、类体系总览

BasicManipuldateProcessor<T>          ← 基类,封装持久化通用逻辑(本文重点)
├── DefaultDataXProcessorManipulate ← DataX 管道的伴生操作基类
│ ├── BatchJobCrontab ← 定时任务配置
│ ├── AddMonitorForEvents ← 告警监控配置
│ └── ExportTISPipelineToDolphinscheduler ← 导出到调度系统
└── OntologyDomainManipulate ← 本体域的伴生操作基类
├── EnableChatBI ← 开启智能问数(需持久化)
└── InferOntologyFromLLM ← LLM 推断本体关系(不持久化)

关键包路径:

路径
BasicManipuldateProcessortis-plugin/.../plugin/manipulate/
DefaultDataXProcessorManipulatetis-plugin/.../datax/
OntologyDomainManipulatetis-plugin/.../plugin/ontology/
EnableChatBI / InferOntologyFromLLMtis-ontology-plugin/.../plugin/ontology/

三、核心接口与基类详解

3.1 IPluginStore.ManipuldateProcessor

这是所有 Manipulate 类必须实现的顶层接口,只有一个方法:

interface ManipuldateProcessor {
void manipuldateProcess(IPluginContext pluginContext,
UploadPluginMeta pluginMeta,
Optional<Context> context);
}

你不需要直接实现这个接口。继承 BasicManipuldateProcessor 后,框架已经帮你实现了。


3.2 BasicManipuldateProcessor<T>

这是所有 Manipulate 插件的基类,封装了完整的持久化流程:

用户触发操作

manipuldateProcess() ← final,不可覆盖

ManipuldateUtils.instance() ← 解析请求上下文(是新增/更新/删除?)

loadPluginStore() ← 抽象方法,子类实现,返回对应的 IPluginStore

[删除] deleteFromStore()
[新增/更新] 判重 → replaceInStore()

afterManipuldateProcess() ← 钩子,子类按需覆盖(如触发同步任务)

你需要实现的方法只有两个:

loadPluginStore()(必须实现)

告诉基类"把数据存到哪里"。从 itemsProcessor.getPluginMeta() 中提取上下文(如 domain 名、pipeline 名),返回对应的 IPluginStore

@Override
protected IPluginStore<MyManipulate> loadPluginStore(
IPluginContext pluginContext,
ManipulateItemsProcessor itemsProcessor) {
// 从 pluginMeta 中取出业务上下文
String domain = itemsProcessor.getPluginMeta()
.getExtraParam(OntologyDomain.NAME_ONTOLOGY_DOMAIN);
KeyedPluginStore.Key<MyManipulate> key =
OntologyDomain.getStoreKey(domain, MyManipulate.class);
return TIS.getPluginStore(key);
}

afterManipuldateProcess()(按需覆盖)

持久化完成后的后续操作。默认是空实现,如果你的操作需要在保存后触发额外逻辑(如同步到 Neo4j、发送通知),在这里写。

@Override
protected void afterManipuldateProcess(
IPluginContext pluginContext,
Optional<Context> context,
ManipulateItemsProcessor itemsProcessor) {
if (itemsProcessor.isDeleteProcess()) {
return; // 删除时通常不需要后续操作
}
// 触发你的业务逻辑
MyService.getInstance().doSomething();
}

3.3 BasicDesc(内部 Descriptor 基类)

每个 Manipulate 插件都需要一个内部静态类 DftDesc(或其他名字),继承对应的 BasicDesc,并加上 @TISExtension 注解。

最关键的方法:isManipulateStorable()

@TISExtension
public static final class DftDesc extends OntologyDomainManipulate.BasicDesc {
@Override
public boolean isManipulateStorable() {
return true; // true = 需要持久化;false = 只执行逻辑,不存储
}
}
isManipulateStorable() 返回值行为
true基类自动完成判重、持久化(replaceInStore/deleteFromStore),然后调用 afterManipuldateProcess
false跳过持久化,直接调用 afterManipuldateProcess

四、开发流程:手把手写一个 Manipulate 插件

以"为本体域开启某个新功能"为例,完整步骤如下。

Step 1:确定你的操作属于哪个核心实体

  • 操作针对 DataX 数据管道 → 继承 DefaultDataXProcessorManipulate
  • 操作针对 本体域 → 继承 OntologyDomainManipulate
  • 操作针对其他实体 → 继承对应的 BasicManipuldateProcessor 子类

Step 2:创建操作类

package com.qlangtech.tis.plugin.ontology;

import com.qlangtech.tis.extension.TISExtension;
import com.qlangtech.tis.plugin.annotation.FormField;
import com.qlangtech.tis.plugin.annotation.FormFieldType;
import com.qlangtech.tis.plugin.annotation.Validator;
import com.qlangtech.tis.plugin.ds.manipulate.ManipulateItemsProcessor;
import com.qlangtech.tis.runtime.module.misc.IControlMsgHandler;
import com.qlangtech.tis.util.IPluginContext;

import java.util.Optional;
import com.alibaba.citrus.turbine.Context;

public class MyNewFeature extends OntologyDomainManipulate {

// 在界面上展示的表单字段
@FormField(ordinal = 1, type = FormFieldType.INPUTTEXT,
validate = {Validator.require})
public String someConfig;

@Override
protected void afterManipuldateProcess(
IPluginContext pluginContext,
Optional<Context> context,
ManipulateItemsProcessor itemsProcessor) {
if (itemsProcessor.isDeleteProcess()) {
return;
}
// 持久化完成后,触发你的业务逻辑
String domain = itemsProcessor.getPluginMeta()
.getExtraParam(OntologyDomain.NAME_ONTOLOGY_DOMAIN);
MyService.activate(domain, this.someConfig);
}

@TISExtension
public static final class DftDesc extends OntologyDomainManipulate.BasicDesc {
public DftDesc() {
super();
}

@Override
public boolean isManipulateStorable() {
return true; // 需要持久化
}

@Override
public String getDisplayName() {
return "My New Feature";
}
}
}

Step 3:判断是否需要持久化

问自己一个问题:用户下次打开界面,需要看到这个操作的配置吗?

  • 需要(如 EnableChatBI 的 LLM 选择)→ isManipulateStorable() 返回 true
  • 不需要(如 InferOntologyFromLLM 只是触发一次推断)→ isManipulateStorable() 返回 false

Step 4:处理删除逻辑

afterManipuldateProcess 在删除操作时也会被调用。通过 itemsProcessor.isDeleteProcess() 判断:

@Override
protected void afterManipuldateProcess(...) {
if (itemsProcessor.isDeleteProcess()) {
// 清理资源,如关闭连接、删除远端配置
MyService.deactivate(domain);
return;
}
// 新增/更新逻辑
MyService.activate(domain, this.someConfig);
}

五、两个典型案例对比

案例 A:EnableChatBI(需持久化 + 有后续操作)

public class EnableChatBI extends OntologyDomainManipulate implements IManipulateStatus {

@FormField(type = FormFieldType.SELECTABLE, ordinal = 1,
validate = {Validator.identity, Validator.require})
public String llm; // 用户选择的大模型

@Override
protected void afterManipuldateProcess(IPluginContext pluginContext,
Optional<Context> context, ManipulateItemsProcessor itemsProcessor) {
if (itemsProcessor.isDeleteProcess()) return;
// 持久化完成后,触发 Neo4j 全量重建
String domainName = itemsProcessor.getPluginMeta()
.getExtraParam(OntologyDomain.NAME_ONTOLOGY_DOMAIN);
OntologySyncQueue.enqueue(
() -> OntologyNeo4jSyncService.getInstance().fullRebuild(domainName));
}

@TISExtension
public static final class DftDesc extends OntologyDomainManipulate.BasicDesc {
@Override
public boolean isManipulateStorable() {
return true; // llm 字段需要持久化,下次打开界面能看到
}
}
}

要点:

  • isManipulateStorable() = truellm 字段会被保存到 XML 配置文件
  • afterManipuldateProcess 里触发异步 Neo4j 同步
  • 实现了 IManipulateStatus 接口,可以在界面上展示当前状态

案例 B:InferOntologyFromLLM(不持久化,只执行一次性任务)

public class InferOntologyFromLLM extends OntologyDomainManipulate {

@FormField(type = FormFieldType.SELECTABLE, ordinal = 1)
public String llm;

@Override
protected void afterManipuldateProcess(IPluginContext pluginContext,
Optional<Context> context, ManipulateItemsProcessor itemsProcessor) {
// 调用 LLM 推断本体关系,结果直接写入本体存储
OntologyPluginMeta ometa = getOntologyPluginMeta(pluginContext, context);
List<OntologyObjectType> objectTypes =
OntologyObjectType.loadAll(ometa.getDomain());
// ... 调用 LLM,解析结果,写入本体 ...
}

@TISExtension
public static final class DftDesc extends OntologyDomainManipulate.BasicDesc {
@Override
public boolean isManipulateStorable() {
return false; // 不需要持久化,每次触发都是全新推断
}
}
}

要点:

  • isManipulateStorable() = false:基类跳过持久化,直接调用 afterManipuldateProcess
  • 所有业务逻辑都在 afterManipuldateProcess

六、常见错误与注意事项

❌ 错误 1:覆盖 manipuldateProcess

// 错误!manipuldateProcess 是 final 的,不能覆盖
@Override
public void manipuldateProcess(...) { ... }

正确做法:把逻辑放到 afterManipuldateProcess


❌ 错误 2:DftDesc 忘记加 @TISExtension

// 错误!没有 @TISExtension,TIS 框架找不到这个 Descriptor
public static final class DftDesc extends OntologyDomainManipulate.BasicDesc { ... }

正确做法:必须加 @TISExtension


❌ 错误 3:DftDesc 没有继承正确的 BasicDesc

// 错误!继承了错误的基类
public static final class DftDesc extends Descriptor<MyFeature> { ... }

正确做法:必须继承对应的 BasicDescOntologyDomainManipulate.BasicDescDefaultDataXProcessorManipulate.BasicDesc)。


❌ 错误 4:在 afterManipuldateProcess 里忘记处理删除分支

// 危险!删除时也会触发 afterManipuldateProcess
@Override
protected void afterManipuldateProcess(...) {
// 如果是删除操作,这里会报错(domain 已不存在)
MyService.activate(domain);
}

正确做法:先判断 itemsProcessor.isDeleteProcess()


⚠️ 注意:loadPluginStore 中的上下文提取

itemsProcessor.getPluginMeta().getExtraParam(key) 是从请求的 pluginMeta 中提取参数。不同的核心实体,参数 key 不同:

核心实体上下文 key获取方式
本体域OntologyDomain.NAME_ONTOLOGY_DOMAINpluginMeta.getExtraParam(...)
DataX 管道DataXNameitemsProcessor.getOriginIdentityId()

七、开发检查清单

开发完成后,对照以下清单自查:

  • 操作类继承了正确的基类(OntologyDomainManipulateDefaultDataXProcessorManipulate
  • 实现了 loadPluginStore()(如果直接继承 BasicManipuldateProcessor
  • 内部类 DftDesc 加了 @TISExtension
  • DftDesc 继承了正确的 BasicDesc
  • isManipulateStorable() 返回值与业务需求一致
  • afterManipuldateProcess 里处理了删除分支
  • 编译通过:mvn compile -q

八、总结

场景做法
需要持久化配置isManipulateStorable() 返回 true,业务逻辑写在 afterManipuldateProcess
只执行一次性任务isManipulateStorable() 返回 false,逻辑全写在 afterManipuldateProcess
持久化后需要触发额外操作覆盖 afterManipuldateProcess,在里面调用业务服务
克隆/重命名场景覆盖 getNewIdentityName() 返回新名称

整个体系的设计思路是模板方法模式:基类 BasicManipuldateProcessor 定义了持久化的完整流程,子类只需填充"存到哪里"(loadPluginStore)和"存完之后做什么"(afterManipuldateProcess)两个空白,其余的校验、判重、存储、删除逻辑由框架统一处理。

有问题随时找团队的同学,欢迎加入 TIS!

· 24 min read

引言

在数据集成平台的开发过程中,我们经常会遇到这样的场景:某个功能的配置项非常多,涉及多个维度的选择,而且这些配置项之间存在依赖关系。比如配置主表与维表的JOIN关系时,需要先选择数据源,再选择目标表,最后设置匹配条件和输出列。如果把所有配置项都堆在一个页面上,不仅用户体验差,而且代码维护困难。

传统的单页面配置方式存在明显的局限性:配置项过多导致页面臃肿、用户容易迷失在复杂的表单中、配置步骤的逻辑关系不清晰。更重要的是,当需要新增类似的复杂配置功能时,往往需要重复大量的开发工作。

TIS(Table Integration System)通过设计一套通用的多步配置插件架构,优雅地解决了这个问题。这套架构不仅让复杂配置变得简单易用,更重要的是,它具有极强的扩展性——新增一个多步配置功能,只需要实现几个接口,编写少量代码即可完成。本文将深入剖析这套架构的设计原理,帮助开发者理解如何构建一个既灵活又易于扩展的配置系统。

多步配置插件的核心架构

整体架构设计

TIS的多步配置插件架构采用了前后端分离的设计思想,通过清晰的接口定义和配置驱动的方式,实现了高度的解耦和可扩展性。

【图片位置1:多步配置插件整体架构图】 说明:展示前端(Angular组件)、后端(Java插件接口)、配置文件(JSON)三者之间的关系,以及数据流转的方向

整个架构可以分为三个核心层次:

  1. 插件宿主层:定义多步配置的容器,负责管理所有步骤
  2. 步骤实现层:每个步骤的具体配置逻辑
  3. 数据解析层:将前端提交的数据转换为后端对象

这三层通过接口进行交互,每一层都可以独立扩展,互不影响。

后端架构设计

后端架构的核心是一套精心设计的接口体系,它定义了多步配置的标准流程和扩展点。

MultiStepsSupportHost:多步配置的宿主

这是多步配置插件的顶层接口,任何需要支持多步配置的插件都需要实现这个接口。它的职责非常明确:

  • 持有所有步骤的配置数据
  • 提供步骤数据的访问接口
  • 定义步骤的元数据(步骤数量、每步的描述等)

通过实现这个接口,一个普通的插件就具备了多步配置的能力。这种设计遵循了"组合优于继承"的原则,让多步配置成为一种可选的能力,而不是强制的约束。

OneStepOfMultiSteps:单步配置的抽象

每个配置步骤都是一个独立的插件,继承自OneStepOfMultiSteps抽象类。这个设计非常巧妙:

  • 每个步骤都是独立的,有自己的配置项和验证逻辑
  • 步骤之间通过接口方法进行数据传递
  • 步骤可以访问前序步骤的配置数据,实现配置的级联

例如,在JOIN配置场景中,第一步选择数据源后,第二步就可以根据选择的数据源加载对应的表列表;第三步又可以根据前两步的选择,展示可用的列信息。这种设计让复杂的配置流程变得清晰可控。

ElementCreatorFactory:数据解析的工厂模式

当用户在前端填写了复杂的表单数据(比如多个匹配条件、多个过滤规则),这些数据需要被解析成后端的Java对象。ElementCreatorFactory就是专门负责这个转换过程的工厂类。

它的设计体现了几个关键思想:

  1. 职责单一:每种数据类型对应一个Factory,只负责解析自己的数据
  2. 配置驱动:通过JSON配置文件指定使用哪个Factory,无需修改代码
  3. 验证集成:在解析过程中同步进行数据验证,及时反馈错误

更重要的是,Factory还负责向前端提供元数据,比如下拉框的选项、字段的约束条件等。这种双向的数据交互机制,让前后端的协作变得非常流畅。

ViewContent:前后端视图类型的桥梁

ViewContent是一个枚举类型,它定义了所有支持的视图类型。每种视图类型对应一种特定的数据结构和UI组件。

这个设计看似简单,实则非常关键:

  • 它是前后端约定的"协议",双方通过ViewContent来识别数据类型
  • 新增一种视图类型,只需要在枚举中添加一个值
  • 前后端可以独立开发,只要遵循ViewContent的约定即可

【图片位置2:前后端数据流转示意图】 说明:展示用户操作 → 前端组件 → ViewContent映射 → ElementCreatorFactory解析 → 后端对象的完整流程

前端架构设计

前端架构同样遵循了组件化和类型驱动的设计原则。

TuplesPropertyType:视图类型的注册机制

TuplesPropertyType是前端对应ViewContent的枚举类型。当后端返回配置数据时,会携带ViewContent信息,前端通过TuplesPropertyType进行映射,决定使用哪个组件来渲染。

这种类型注册机制的优势在于:

  • 新增视图类型时,只需要在switch语句中添加一个case分支
  • 每个组件都是独立的,可以单独开发和测试
  • 类型安全,编译期就能发现错误

组件化的步骤UI

每个配置步骤对应一个Angular组件,这些组件都继承自BasicTuplesViewComponent基类。基类提供了通用的功能:

  • 数据绑定和验证
  • 错误信息展示
  • 与父组件的通信

具体的步骤组件只需要关注自己的业务逻辑,比如渲染表单、处理用户交互等。这种设计让组件的复用性非常高。

数据双向绑定与验证

前端使用Angular的双向绑定机制,用户的每次输入都会实时同步到数据模型中。当用户提交配置时,前端会将数据序列化为JSON格式,发送给后端。

验证逻辑分为两层:

  1. 前端验证:即时反馈,提升用户体验
  2. 后端验证:保证数据安全,防止恶意提交

后端验证失败时,错误信息会通过统一的格式返回给前端,前端根据字段路径将错误信息显示在对应的表单项旁边。

架构的扩展性设计

这套架构最大的价值在于它的扩展性。当需要新增一个多步配置功能时,开发工作量非常小,这得益于几个关键的设计决策。

接口抽象的威力

通过定义清晰的接口,架构将"做什么"和"怎么做"完全分离。

以JOIN配置为例,当我们需要新增"过滤条件"功能时,只需要:

  1. 创建一个新的数据类TableJoinFilterCondition,实现IMultiElement接口
  2. 创建对应的Factory类,实现ElementCreatorFactory接口
  3. 在ViewContent枚举中添加一个新值

这三步完成后,后端的扩展就完成了。整个过程不需要修改任何现有代码,完全符合"开闭原则"——对扩展开放,对修改关闭。

接口抽象的另一个好处是,它强制开发者思考功能的本质。比如ElementCreatorFactory接口定义了几个关键方法:

  • parsePostMCols:解析前端提交的数据
  • appendExternalJsonProp:向前端提供元数据
  • getViewContentType:返回视图类型

任何需要解析复杂表单数据的场景,都可以套用这个模式。开发者不需要重新设计数据流转的逻辑,只需要实现这几个方法即可。

配置驱动的灵活性

JSON配置文件在这套架构中扮演了"粘合剂"的角色。它连接了Java字段和ElementCreatorFactory,让两者可以独立变化。

配置文件的结构非常简洁:

{
"字段名": {
"label": "显示标签",
"elementCreator": "Factory类的全限定名",
"enum": "获取初始数据的方法",
"viewtype": "视图类型",
"help": "帮助文本"
}
}

这种配置驱动的设计带来了几个好处:

  1. 热插拔:修改配置文件即可改变行为,无需重新编译
  2. 可视化:配置文件清晰地展示了字段和Factory的对应关系
  3. 解耦:Java代码和Factory类之间没有直接依赖

更重要的是,这种设计让非开发人员也能理解系统的配置逻辑。当需要调整字段顺序、修改标签文本时,只需要编辑JSON文件即可。

组件化的可复用性

前端的组件化设计让UI的复用变得非常简单。

每个组件都是自包含的,它有自己的模板、样式和逻辑。当需要创建一个新的配置步骤时,可以参考现有组件的实现,快速搭建出新的UI。

组件之间通过标准的Input和Output进行通信:

  • Input接收父组件传递的数据和配置
  • Output向父组件发送事件和数据变更

这种标准化的通信方式,让组件可以像乐高积木一样自由组合。比如,一个下拉选择组件可以在多个不同的配置步骤中使用,只需要传入不同的选项数据即可。

Angular的依赖注入机制进一步增强了组件的可测试性。每个组件可以独立进行单元测试,不需要启动整个应用。

实际应用:主表与维表JOIN配置

让我们通过一个实际的应用场景,来看看这套架构是如何工作的。

JOIN配置的业务场景

在数据集成场景中,经常需要将主表和维表进行JOIN操作,以实现数据的宽表化。比如订单表需要关联用户表,获取用户的详细信息;商品表需要关联分类表,获取分类名称等。

这个配置过程涉及多个步骤:

  1. 选择维表的数据源(可能有多个数据库)
  2. 选择具体的维表
  3. 设置JOIN的匹配条件(主表的哪个字段等于维表的哪个字段)
  4. 设置过滤条件(可选,比如只关联有效的记录)
  5. 选择需要输出的列
  6. 配置缓存策略

如果用传统的单页面方式,这个表单会非常复杂。而使用多步配置插件,可以将这个过程分解为三个清晰的步骤。

【图片位置3:JOIN配置三步向导界面截图】 说明:展示三个步骤的界面,第一步选择数据源,第二步选择表,第三步配置匹配条件、过滤条件和输出列

三步配置的设计思路

第一步:选择数据源

这一步的职责非常明确:让用户从已配置的数据源中选择一个。实现上,创建一个JoinerSelectDataSource类,继承OneStepOfMultiSteps,定义一个dbName字段用于存储选择的数据源名称。

前端对应一个简单的下拉选择组件,选项从后端的数据源列表中获取。

第二步:选择目标表

这一步依赖第一步的选择结果。JoinerSelectTable类可以通过getOneStepOf方法获取第一步的配置,从而知道用户选择了哪个数据源,进而加载该数据源下的表列表。

这种步骤间的数据传递机制,让配置流程具有了"智能"——后续步骤可以根据前面的选择动态调整。

第三步:配置匹配条件和输出列

这是最复杂的一步,需要配置多个匹配条件、多个过滤条件,还要选择输出列。

JoinerSetMatchConditionAndCols类定义了多个字段:

  • matchCondition:匹配条件列表
  • filterConditions:过滤条件列表
  • targetCols:输出列列表
  • colPrefix:列名前缀
  • cache:缓存配置

每个字段都通过@FormField注解标注,指定了字段类型、验证规则等。关键的是,matchCondition和filterConditions字段的类型是MULTI_SELECTABLE,表示这是一个复杂的多元组配置。

过滤条件功能的扩展实现

最初的JOIN配置只支持匹配条件,后来需要新增过滤条件功能。这个扩展过程完美展示了架构的灵活性。

后端扩展

  1. 创建TableJoinFilterCondition类,定义5个字段:tableType(主表/维表)、columnName(列名)、operator(运算符)、valueType(值类型)、value(值)
  2. 创建TableJoinFilterConditionCreatorFactory类,实现数据解析逻辑,并提供运算符、值类型等元数据
  3. 在ViewContent枚举中添加TableJoinFilterCondition值
  4. 在JoinerSetMatchConditionAndCols类中添加filterConditions字段
  5. 在JSON配置文件中添加filterConditions的配置项

前端扩展

  1. 创建TableJoinFilterConditionComponent组件,渲染过滤条件的表单
  2. 在TuplesPropertyType枚举中添加TableJoinFilterCondition值
  3. 在类型映射的switch语句中添加对应的case分支
  4. 在Angular模块中声明新组件

整个扩展过程,没有修改任何现有的代码,只是新增了文件和配置。这就是良好架构设计带来的红利。

技术优势总结

开发效率的显著提升

通过这套架构,新增一个多步配置功能的开发成本大幅降低。

传统方式下,开发一个复杂的配置功能可能需要:

  • 设计数据库表结构
  • 编写后端的CRUD接口
  • 设计前端的表单布局
  • 实现前后端的数据交互
  • 编写验证逻辑
  • 处理各种边界情况

整个过程可能需要几天甚至一周的时间。

而使用多步配置插件架构,开发者只需要:

  • 定义数据模型(实现接口)
  • 实现Factory类(套用模板)
  • 编写前端组件(参考现有组件)
  • 添加配置项(编辑JSON文件)

熟练的开发者可以在半天内完成一个新功能的开发。这种效率的提升,来自于架构的标准化和模板化。

用户体验的优化

向导式的配置流程,让用户不再面对复杂的表单,而是一步步地完成配置。每一步的目标都很明确,用户不会感到迷茫。

步骤之间的数据联动,让配置过程更加智能。比如选择了数据源后,表列表会自动加载;选择了表后,列信息会自动展示。用户不需要手动刷新或重新加载。

实时的验证反馈,让用户在填写过程中就能发现错误,而不是提交后才被告知配置有误。这种即时反馈大大提升了配置的成功率。

系统可维护性的增强

清晰的架构边界,让代码的职责非常明确。当需要修改某个功能时,开发者可以快速定位到对应的类和组件,不会影响其他功能。

接口驱动的设计,让系统具有很好的可测试性。每个接口的实现都可以独立测试,不需要启动整个系统。这让单元测试变得简单可行。

配置文件的使用,让系统的行为更加透明。运维人员可以通过查看配置文件,了解系统的配置逻辑,而不需要阅读代码。

扩展场景展望

这套多步配置插件架构的应用场景远不止JOIN配置。它可以应用于任何需要复杂配置的场景。

数据转换规则配置:数据同步过程中,经常需要对数据进行转换,比如字段映射、类型转换、格式化等。这些转换规则的配置可以分为多个步骤:选择源字段、选择转换函数、配置参数、预览结果。

复杂数据源配置:某些数据源的配置非常复杂,比如Kafka需要配置集群地址、Topic、消费组、序列化方式等。通过多步配置,可以将这些配置项分组,让用户逐步完成。

工作流编排配置:数据处理流程可能包含多个节点,每个节点有自己的配置。通过多步配置,可以让用户先定义流程结构,再逐个配置节点参数。

机器学习模型配置:训练一个机器学习模型需要配置数据源、特征工程、算法参数、评估指标等。这个过程天然适合用多步配置来实现。

这些场景的共同特点是:配置项多、步骤间有依赖、需要向导式的用户体验。只要符合这些特点,都可以使用这套架构来实现。

总结

TIS的多步配置插件架构,通过精心设计的接口体系、配置驱动的灵活机制、组件化的前端实现,构建了一个既强大又易于扩展的配置系统。

这套架构的核心思想可以总结为几点:

  1. 接口抽象:通过接口定义标准流程,让扩展变得简单
  2. 配置驱动:用配置文件连接各个组件,实现解耦
  3. 组件化:前后端都采用组件化设计,提高复用性
  4. 类型安全:通过枚举类型实现前后端的类型映射,减少错误

对于开发者来说,这套架构提供了一个很好的参考:当面对复杂的配置场景时,不要试图用一个大而全的表单来解决,而是应该思考如何将配置过程分解为多个步骤,如何设计接口让系统具有扩展性,如何通过配置来驱动行为。

良好的架构设计,不仅能提升开发效率,更能让系统具有持续演进的能力。TIS的多步配置插件架构,正是这样一个经过实践检验的优秀设计。


本文基于TIS(Table Integration System)开源项目的实际开发经验总结而成。TIS是一个企业级的数据集成平台,致力于为数据工程师提供高效、易用的数据同步和处理工具。

· 19 min read
百岁

引言

在数字化转型的今天,企业越来越依赖实时数据处理来支撑业务决策。然而,实时任务一旦出现故障,可能导致数据丢失、业务中断等严重后果。传统的做法是安排专人盯盘监控,这不仅效率低下,而且无法做到24小时不间断监控。

TIS平台最新上线的Flink任务智能监控与报警系统,彻底解决了这一痛点。系统能够自动发现任务异常,并第一时间通过钉钉、企业微信、邮件等多种方式通知相关人员,真正做到"任务有问题,立即能知道"。

以下是钉钉群中接收到一条TIS发送的报警消息:

传统监控方式的痛点

在引入自动化监控之前,企业通常采用以下方式来监控Flink任务:

人工定期检查

运维人员需要定期登录Flink管理界面,逐个检查任务状态。这种方式存在明显问题:

  • 响应滞后: 任务失败后可能几小时才被发现
  • 人力成本高: 需要专人盯盘,占用大量人力
  • 覆盖不全: 夜间、周末、节假日可能无人值守
  • 容易遗漏: 任务数量多时,容易漏检

简单的脚本监控

一些技术团队会编写简单的Shell脚本定时检查任务状态,但这种方式也有局限:

  • 通知渠道单一: 通常只能发邮件,不够及时
  • 维护成本高: 每增加一个任务就要修改脚本
  • 缺乏统一管理: 脚本分散在各处,难以维护
  • 报警信息简陋: 只能提供基本的状态信息,缺少详细数据

第三方监控工具

使用Prometheus、Grafana等监控工具,需要:

  • 学习成本高: 需要掌握新的技术栈
  • 配置复杂: 需要编写大量配置文件
  • 集成困难: 与现有系统集成需要开发工作
  • 额外成本: 需要维护额外的监控系统

TIS智能监控方案的优势

面对上述痛点,TIS平台推出了开箱即用的智能监控方案,具有以下显著优势:

1. 零配置,开箱即用

TIS的监控功能是平台内置的,无需安装额外的软件或服务。系统启动后,监控功能自动运行:

  • 自动发现: 自动监控所有运行中的Flink任务,无需手动配置
  • 实时响应: 每5秒检测一次,任务异常后最快5秒即可发现
  • 智能判断: 区分正常停止和异常中断,避免误报

以下是Flink任务监控系统工作流程:

流程说明:

  1. 监控启动: TIS平台启动时自动启动FlinkJobsMonitor监控器
  2. 定时轮询: 每5秒执行一次任务状态检测
  3. 状态检测: 获取所有Flink任务的当前状态
  4. 状态比对: 与上次记录的状态进行对比,判断是否发生变化
  5. 异常判断: 识别以下异常情况:
    • 任务失败 (RUNNING → FAILED)
    • 任务丢失 (RUNNING → LOST)
    • 异常取消 (RUNNING → CANCELED 且非手动停止)
  6. 报警触发: 构建包含任务详细信息的AlertTemplate对象
  7. 多渠道发送: 并发发送到所有已配置的报警渠道
  8. 模板渲染: 使用Velocity模板渲染个性化的报警消息
  9. 消息发送: 根据不同渠道的协议发送报警通知
  10. 状态更新: 更新任务状态记录,继续下一轮监控

2. 多渠道报警,信息必达

TIS支持5种主流的报警渠道,满足不同场景的需求:

📧 邮件报警

  • 适合重要告警的归档和追溯
  • 支持发送到多个邮箱
  • 提供精美的HTML格式报告

💬 钉钉报警

  • 研发团队最常用的即时通讯工具
  • 消息直达钉钉群,响应迅速
  • 支持@指定人员,确保信息触达

🏢 企业微信报警

  • 适合使用企业微信的团队
  • 支持Markdown格式,信息清晰
  • 可@多人协同处理

🚀 飞书报警

  • 字节跳动旗下的协作平台
  • 支持丰富的卡片样式
  • 信息展示更加美观

🔗 HTTP回调

  • 最灵活的集成方式
  • 可对接任何支持HTTP接口的系统
  • 适合与现有运维系统集成

核心优势: 可以同时配置多个渠道!例如:邮件用于归档,钉钉用于即时响应,企业微信通知运维团队,真正做到"多管齐下,确保必达"。

3. 报警信息丰富,一目了然

TIS的报警消息不是简单的"任务失败"几个字,而是提供了全面的任务信息:

基础信息

  • 任务名称和状态(失败/丢失/异常取消)
  • 开始时间、结束时间、运行时长
  • 一键跳转到Flink Web UI查看详情

智能分析

  • 自动计算运行时长
  • 标注异常类型(首次失败/重复失败)
  • 提供历史状态对比

以下是通过邮件渠道接收到一封TIS报警消息:

举个例子: 凌晨3点,某个数据同步任务突然失败。5秒后,值班人员的钉钉群就收到了报警消息,上面清楚地显示:

  • 哪个任务出问题了
  • 什么时候出的问题
  • 已经运行了多长时间
  • 点击链接可以直接查看Flink控制台 值班人员不用登录系统,在手机上就能了解全部情况,快速做出响应。

4. 灵活定制,满足个性化需求

虽然TIS提供了开箱即用的默认配置,但也充分考虑了不同企业的个性化需求:

自定义报警内容

  • 支持Velocity模板,可以自定义报警消息的格式和内容
  • 可以添加企业特有的信息(如项目名称、负责人联系方式等)
  • 不同渠道可以使用不同的消息格式

灵活的配置管理

  • Web界面配置,无需修改代码
  • 支持配置多个同类型渠道(如:生产环境钉钉群、测试环境钉钉群)
  • 配置立即生效,无需重启服务

安全机制

  • 钉钉和飞书支持密钥签名,防止Webhook泄露
  • 邮件支持SSL加密传输
  • HTTP回调支持自定义请求头,可配置认证Token

5. 稳定可靠,久经考验

TIS的监控系统在设计上充分考虑了稳定性和可靠性:

异常隔离

  • 单个任务检测失败不影响其他任务
  • 某个报警渠道发送失败不影响其他渠道
  • 监控系统异常不影响Flink任务本身的运行

性能优化

  • 轻量级设计,对系统性能影响极小
  • 高效的状态管理,避免重复报警
  • 合理的轮询频率,在实时性和性能间取得平衡

久经考验

  • 设计理念借鉴了Apache StreamPark等成熟开源项目
  • 采用TIS平台久经考验的插件化架构
  • 已在多个生产环境稳定运行

真实使用场景

场景一:数据库增量同步任务监控

某电商企业使用TIS搭建了MySQL到Elasticsearch的实时数据同步管道,用于实时更新商品搜索索引。

面临的挑战:

  • 有20多个数据库表需要同步
  • 任务需要7×24小时不间断运行
  • 一旦中断会影响用户搜索体验

使用TIS监控后:

  • 配置了钉钉报警和邮件报警双通道
  • 某次因为网络抖动导致一个同步任务失败
  • 5秒内研发群收到钉钉报警,10分钟内问题得到解决
  • 同时邮件报警提供了完整的事件记录,便于后续分析

效果:

  • 故障响应时间从平均2小时降低到10分钟
  • 再也不用安排人员夜间值班
  • 运维成本降低60%

场景二:大数据ETL任务监控

某金融企业使用TIS处理每日的交易数据,将数据从业务库同步到数据仓库。

面临的挑战:

  • 涉及多个业务系统,数据量大
  • 对数据时效性要求高
  • 需要及时发现并处理异常

使用TIS监控后:

  • 配置了企业微信报警,消息直达运维群
  • 可以@不同的负责人处理不同任务的异常
  • 报警消息中包含详细的任务信息和性能指标
  • 通过HTTP回调将报警信息集成到公司的运维平台

效果:

  • 异常任务发现率100%
  • 平均处理时间减少70%
  • 建立了完整的运维响应机制

如何使用TIS监控功能

第一步:配置报警渠道

以配置钉钉报警为例,整个过程只需3分钟:

1. 创建钉钉群机器人

  • 打开需要接收报警的钉钉群
  • 点击"群设置" → "智能群助手" → "添加机器人" → "自定义"
  • 设置机器人名称,安全设置选择"加签"
  • 复制Webhook URL和密钥

2. 在TIS中配置报警渠道

  • 登录TIS平台
  • 进入"全局配置" → "报警渠道"
  • 点击"添加",选择"DingTalk"
  • 填写配置:
    • 渠道名称:如"生产环境钉钉"
    • Webhook URL:粘贴刚才复制的URL
    • 密钥:粘贴刚才复制的密钥
    • (可选)填写需要@的人员手机号
  • 保存配置

3. 完成

  • 配置保存后立即生效
  • 系统自动开始监控所有Flink任务
  • 有异常时自动发送报警到钉钉群

第二步:测试报警功能

配置完成后,建议先测试一下:

  1. 在TIS中找一个测试任务
  2. 手动停止该任务
  3. 等待5-10秒
  4. 检查钉钉群是否收到报警消息

收到报警说明配置成功!消息中会显示:

  • 任务名称和状态
  • 运行时长
  • 查看详情的链接

第三步:配置多个渠道(可选)

根据需要,您还可以配置其他报警渠道:

  • 邮件: 用于重要告警的归档,方便后续查询
  • 企业微信: 通知运维团队
  • 飞书: 如果团队使用飞书办公
  • HTTP回调: 集成到现有的运维系统

每个渠道的配置方式类似,都是在"报警渠道"页面点击"添加",填写相应的配置信息即可。

核心价值总结

TIS Flink任务智能监控系统为企业带来的核心价值:

💰 降低成本

  • 不需要专人盯盘,节省人力成本
  • 无需部署额外的监控系统,降低IT成本
  • 快速响应减少故障损失

⚡ 提升效率

  • 5秒发现异常,响应速度提升数十倍
  • 自动化监控,释放运维人力
  • 多渠道通知,确保信息必达

🛡️ 保障稳定

  • 7×24小时不间断监控
  • 智能判断,避免误报和漏报
  • 详细信息,快速定位问题

📊 易于管理

  • Web界面配置,无需技术背景
  • 统一管理所有Flink任务
  • 支持多团队、多项目使用

🔧 灵活扩展

  • 5种主流报警渠道,满足不同需求
  • 支持自定义报警模板
  • 可集成到现有运维体系

技术实现亮点

虽然本文面向非技术人员,但有必要简单介绍一下TIS监控系统的技术优势,这些设计保证了系统的稳定可靠:

借鉴成熟开源项目

  • 参考Apache StreamPark的监控设计理念
  • 站在开源社区的肩膀上,避免重复造轮车
  • 结合TIS平台特点进行优化

插件化架构

  • 新增报警渠道无需修改核心代码
  • 各个报警渠道独立运行,互不影响
  • 方便后续扩展更多渠道

智能状态管理

  • 记录任务的历史状态,避免重复报警
  • 自动清理过期数据,防止内存泄漏
  • 高性能的状态比对算法

模板化设计

  • 统一的消息模板变量体系
  • 用户可以自由定制报警内容
  • 同一套模板适用于所有渠道

完善的安全机制

  • 支持密钥签名验证(钉钉、飞书)
  • 支持SSL加密传输(邮件)
  • 支持自定义认证(HTTP回调)

结语

在数字化时代,数据就是企业的生命线。实时数据处理任务的稳定运行,直接关系到业务的正常开展。TIS Flink任务智能监控系统,让企业可以放心地运行大规模实时数据管道,再也不用担心任务异常无人知晓。

开箱即用的设计理念、多渠道的报警支持、丰富详细的报警信息,让TIS的监控功能成为运维人员的得力助手。无论是中小企业还是大型互联网公司,都能从中受益。

如果您正在使用TIS平台,现在就可以开始配置报警渠道,让监控功能为您的Flink任务保驾护航。如果还没有使用TIS,欢迎访问我们的官网了解更多信息。

· 9 min read
百岁

前言

说到依赖反转原则百度百科是这样介绍的,如下:

在面向对象编程领域中,依赖反转原则(Dependency inversion principle,DIP)是指一种特定的解耦(传统的依赖关系创建在高层次上,而具体的策略设置则应用在低层次的模块上)形式,使得高层次的模块不依赖于低层次的模块的实现细节,依赖关系被颠倒(反转),从而使得低层次模块依赖于高层次模块的需求抽象。

大家应该会马上会联想到Spring Framework,在介绍Spring Framework框架常常会提及依赖反转原则,看上面的介绍估计会云里雾里,说得通俗一点,该原则的初衷要求服务提供者与服务调用者在代码实现层面实现解耦

为了加深理解,经常会提到的一个例子,以前古时候的包办婚姻,假如是男方到了适婚年龄,只要把自己的条件、要求告诉媒婆。接下来找合适对象的过程就交给媒婆就行了,男方只需要负责到时候入洞房就行了。对于男方来说, 婚姻和他是服务提供者和消费者的关系,由于引入了媒婆的角色,男方省去了谈婚论嫁的麻烦过程,只需要专注于核心业务-入洞房。最终使得整个过程显得简单高效,形式格外优雅。

因此,依赖反转原则DIP作为一个朴素的原则存在,可以应用到软件设计领域每一个流程环节当中,而不仅仅适用于Spring Framework当中。

本文就以依赖反转原则DIP在TIS增量实时数据通道的设计、实现过程中如何利用这一原则来优化设计、实现流程进行阐述。

实现实时增量数据管道需求

在TIS中为用户提供了基于Flink端到端的实时增量数据通道功能,市面上已经提供了基于Flink和Flink-CDC的实时流同步工具,从用户反馈来看已经很方便了,那为什么还要通过TIS来 使用Flink-CDC呢?

这是一个非常好的问题,要回答这个问题,首先我们需要从用户的角度了解用户到底需要什么?然后从需求出发设计并且构建出用户体验达到极致的产品。

大数据流计算领域,用户的核心需求是:

  1. 可追溯操作历史的控制系统,这样可以方便回滚历史操作。
  2. 不关心算子实现细节,流计算的使用者往往是对Flink不了解的数据分析人员,所以在产品使用体验上需要屏蔽底层技术细节。
  3. 可扩展的端类型:Flink-CDC从3.0版本支持的Connectors,只支持了有限个数的基于增量监听CDC技术的Source端 ,和少量Sink端实现,如:Doris和StarRocks的Sink端类型。还远远没有达到用户实际生产场景下的端类型。所以,需要提供在更高层次上,通过便捷方式扩展Source和Sink端类型的手段。

TIS正式为了弥补以上三个使用Flink-CDC框架中的不足而开发的。

具体实现

下面具体对以上第2点进行进行说明,配置并且触发执行基于Flink-CDC的数据管道具体通过以下步骤完成

编辑

构建DataStreamSource步骤中,通过调用Flink-CDC提供的API代码,可以方便订阅到如MySQL的增量更新消息,如下代码:

public static void main(String[] args) throws Exception {
MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
.hostname("yourHostname")
.port(yourPort)
.databaseList("yourDatabaseName") // set captured database, If you need to synchronize the whole database, Please set tableList to ".*".
.tableList("yourDatabaseName.yourTableName") // set captured table
.username("yourUsername")
.password("yourPassword")
.deserializer(new JsonDebeziumDeserializationSchema()) // converts SourceRecord to JSON String
.build();

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// enable checkpoint
env.enableCheckpointing(3000);

env
.fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL Source")
// set 4 parallel source tasks
.setParallelism(4)
.print().setParallelism(1); // use parallelism 1 for sink to keep message ordering

env.execute("Print MySQL Snapshot + Binlog");
}

通过以上流计算的流程中可以使用创建出MySqlSource<String> mySqlSourceSource加入到各种算子中去进行计算。

使用SQL的方式将Stream Source 注册为Flink Table:

CREATE TABLE mysql_source (...) WITH (
'connector' = 'mysql-cdc',
'scan.startup.mode' = 'earliest-offset', -- Start from earliest offset
'scan.startup.mode' = 'latest-offset', -- Start from latest offset
'scan.startup.mode' = 'specific-offset', -- Start from specific offset
'scan.startup.mode' = 'timestamp', -- Start from timestamp
'scan.startup.specific-offset.file' = 'mysql-bin.000003', -- Binlog filename under specific offset startup mode
'scan.startup.specific-offset.pos' = '4', -- Binlog position under specific offset mode
'scan.startup.specific-offset.gtid-set' = '24DA167-0C0C-11E8-8442-00059A3C7B00:1-19', -- GTID set under specific offset startup mode
'scan.startup.timestamp-millis' = '1667232000000' -- Timestamp under timestamp startup mode
...
)

以上是Flink-CDC提供的标准化的Demo案例。

在这里我们重新用依赖反转原则来思考一下,这个构建流程是否有违背该原则?确实,从用户的角度来说,用户只关心最终构建出来的MySqlSource<String>实例,至于构建该实例的过程用户并不关心,所以在设计过程需要将 MySqlSource<String>实例构建过程与它的调用者之间进行解耦合。

是时候发挥TIS的作用了,TIS需要发挥实例容器的作用,由TIS根据用户配置的Source端参数自动地创建MySqlSource<String>实例, 在运行时自动注入到执行流程中。

  • 配置Source/Sink Connector
  • 直接引用TIS注入的SourceStream实例
  • 当用户选择Flink SQL类型脚本,直接引用已经注册完成的Table名即可

以上具体提供注入实例的封装工厂是:

以上两段代码的执行逻辑类似Spring FactoryBean 执行逻辑,实现容器预定的扩展工厂接口,运行期由容器负责初始化,继而将实例注入到需要反向依赖的实例中。

总结

本文介绍了利用依赖反转原则在TIS中实现实时增量通道的优化方法,可使最终用户最大限度地关注流式计算核心业务本身,其他的琐碎的与实例初始化相关的工作都交给TIS来完成即可。

与此类似的功能优化,在TIS实现过程中还有很多,会在日后的博客分享中陆续发表。

· 8 min read
百岁

前言

MongoDB是一个基于分布式文件存储的开源数据库系统,其内容存储格式为BSON(一种类json的二进制形式),其主要特点是高性能、易部署、易使用,存储数据支持高度事务性且支持完全索引,包括地理空间索引、散列索引和全文索引,还有一个比较大的优势是其适应于海量数据的存储,其数据被分散在不同的服务器上,以自动分区数据。

MongoDB在实际生产环境中有很多应用场景,利用器BSON数据结构可以在运行期动态扩展数据Schema。

MongoDB被TIS整合进了数据集成方案,通过TIS中可以方便对MongoDB端进行读取或者写入操作,方便实现对MongoDB的数据迁移、实时容灾备份、异构数据端(如:Doris)实时同步实现复杂OLAP操作。

本文就实际操作过程中,发现从MongoDB中不能预先读取表Schema,对此进行了优化,并且对此优化过程作以详细介绍。

发现瓶颈

MongoDB作为文档型数据库的代表,区别于传统关系型数据库MySQL,MongoDB的表数据结构Schema在运行期是可变的,而不是像MySQL那样通过Create Table DDL预先定义好表Schema。因此,在为MongoDB做数据集成操作时带来一个麻烦事儿,需要通过手工配置的方式 为读取MongoDB的表作为依据。例如,用户通过Alibaba DataX来读取MongoDB需编写DataX Reader任务配置:

https://github.com/alibaba/DataX/blob/master/mongodbreader/doc/mongodbreader.md
  {
"job": {
"content": [
{
"reader": {
"name": "mongodbreader",
"parameter": {
"address": ["127.0.0.1:27017"],
"dbName": "tag_per_data",
"collectionName": "tag_data12",
"column": [
{
"name": "unique_id",
"type": "string"
},
{
"name": "sid",
"type": "string"
},
{
"name": "user_id",
"type": "string"
},
{
"name": "auction_id",
"type": "string"
},
{
"name": "content_type",
"type": "string"
},
{
"name": "pool_type",
"type": "string"
}
]
}
},
"writer": {
"name": "odpswriter"
}
}
]
}
}

以上配置文件中 reader mongodbreader需要配置对应表的列枚举信息,列中存在BSON类型的列,还存在拆列的问题,会更加复杂,配置过程虽然简单,但还是很容易会出错,特别在配置type属性时。

优化

思路

优化思路,是否可以通过MongoDB的JDBC客户端,通过反射的方式得到表的列信息列表。然后通过模版机制(如:velocity)自动生成DataX配置文件中column配置。

尝试通过MongoDB Client API读取表Schema元数据信息,我们可以尝试从collection中读取一条记录,然后通过解析记录获得Schema记录,如下:

  var schemaObj = db.users.findOne();

遍历记录的所有列

  void printSchema(obj) {
for (var key in obj) {
print(indent, key, typeof obj[key]) ;
}
};

可以将Schema打印出,如下:

这非常酷,而用户自定义Collection往往存在子属性,希望将这些子属性进行拆接打平,可以导入到下游目标端中。

我们可以优化以上代码:

function printSchema(obj, indent) {
for (var key in obj) {
print(indent, key, typeof obj[key]) ;
if (typeof obj[key] == "object") {
printSchema(obj[key], indent + "\t")
}
}
};
printSchema(schemaObj,"");

重新执行以上代码,将会打印:

1_nWI-bWoJHWE1WU2Mgk8z-A.webp

在TIS中具体实现

TIS中读取MongoDB表Schema实现方式沿袭以上思路,另外添加额外的工序:

  1. 在控制台中设置尝试读取的记录数,由于MongoDB Collection中每条记录的列数量和类型不一定相同的,可以尝试读取多条Collection中记录,将每条记录Schema进行Merge最终获得Schema。
  2. 读取每条记录类型为BsonType.DOCUMENT的列类型,将内部子列与父列通过'.'号就行连接,打平成为新的列,例如:"user.name","user.age"
  3. 将获得到的Schema在前台展示,用户可以通过在表单中对Schema结构进行微调,以确认最终的Schema结构。

以下为 MongoColumnMetaData.java GitHub中路径代码中的片段:

/tis-datax-mongodb-plugin/src/main/java/com/qlangtech/tis/plugin/datax/mongo/MongoColumnMetaData.java#L69
    public static void parseMongoDocTypes(boolean parseChildDoc, List<String> parentKeys //
, Map<String, MongoColumnMetaData> colsSchema, BsonDocument bdoc) {
int index = 0;
BsonValue val;
String key;
MongoColumnMetaData colMeta;
List<String> keys = null;

for (Map.Entry<String, BsonValue> entry : bdoc.entrySet()) {
val = entry.getValue();
keys = ListUtils.union(parentKeys, Collections.singletonList(entry.getKey()));
key = String.join(MongoCMeta.KEY_MONOG_NEST_PROP_SEPERATOR, keys);
colMeta = colsSchema.get(key);
if (colMeta == null) {
colMeta = new MongoColumnMetaData(index++, key, val.getBsonType(), 0,
(val.getBsonType() == BsonType.OBJECT_ID));
colsSchema.put(key, colMeta);
} else {
if (colMeta.getMongoFieldType() != BsonType.STRING //
&& !val.isNull() && colMeta.getMongoFieldType() != val.getBsonType()) {
//TODO: 前后两次类型不同
// 则直接将类型改成String类型
colMeta = new MongoColumnMetaData(index++, key, BsonType.STRING);
colsSchema.put(key, colMeta);
}
}
if (!val.isNull()) {
if (colMeta.getMongoFieldType() == BsonType.DOCUMENT && val.isDocument()) {
parseMongoDocTypes(true, keys, parseChildDoc ? colsSchema : colMeta.docTypeFieldEnum, val.asDocument());
}
if (colMeta.getMongoFieldType() == BsonType.STRING) {
colMeta.setMaxStrLength(val.asString().getValue().length());
}
colMeta.incrContainValCount();
}

}
}

MongoDB Reader 页面设置预读记录数,尝试读取Collection记录条数以分析出Collection的Schema,为确保最终Schema准确,可以适当将预读记录数设置得大一些。

通过TIS解析得到的Schema会在以下页面中展示结果,用户可进行微调,确定是否需要该列,变更字段类型等。

通过以上优化,最终完成数据同步通道定义,可以有效避免用户手动输入MongoDB表 Schema,达到了最大限度地降低了出错概率,且提高了工作效率

总结

数据集成领域,有大部分的端类型是和MongoDB类似属于 Schemaless的,例如Kafka,Redis,基于文件的Hdfs,FTP等等。不像MySQL这样的具有明确预定义Schema的数据源可以通过读取MetaData的方式得到Schema,从而进行自动化操作。

TIS的初衷是构建一款高度傻瓜化的DataOps数据集成软件,面向一线非技术人员,他们精通业务,在具体操作过程中不需要了解具体MongoDB中的字段类型,有哪些字段。 整个操作流程,只需要轻点鼠标,TIS会帮助用户自动生成所需配置。

借鉴这种操作思路,可以扩展到其他Schemaless的数据端读取流程上,例如Kafka,Redis,基于文件的Hdfs,FTP,将会极大地提高执行数据集成的效率。

· 3 min read
百岁

前言

在TIS 4.0.0 版本主要的功能是将原先但节点运行的组件扩展到分布式云环境中,以下类图中有三个组件需要依赖到ServerPortExport组件,

  1. Kubernete Powerjob Server
  2. Kubernete Flink Session
  3. Kubernete Flink Application

ServerPortExport 组件负责在K8S组件(ReplicaSet)发布过程中将目标端口以不同的方式发布(Ingress,LoadBalance,NodePort)

编辑

遇到问题

ServerPortExport 组件聚合到不同的组件中,在具体运行过程中需要根据聚合类的不同有不同的初始值,

例如,聚合在K8SDataXPowerJobServer类中初始值为7700,而当聚合在BasicFlinkK8SClusterCfg中的初始值为8081,当然,直观来说,最简单的办法是根据聚合到不同的类,创建不同的ServerPortExport的子类从而来设置不同的初始值, 但这会创建出大量的冗余代码,所以,并不可取。

解决办法

在运行期,根据所在聚合类的Descriptor来动态设置 ServerPortExport.serverPort 属性的值

编辑

具体需要做以下功能:

  1. 创建 DefaultExportPortProvider接口,get方法返回对应的端口默认值
  2. BasicFlinkK8SClusterCfgK8SDataXPowerJobServer对应的 Descriptor分别实现以上接口
  3. 在运行期将Descriptor序列化成Json步骤中,需要将Descriptor实例与当前运行的线程绑定,这部分功能在Json序列化过程中执行,为此需要添加新类DescriptorsJSONResult
  4. 为类DescriptorsJSONResult注册到Json序列化注册器中
    JsonUtil.java
       ObjectWriter descSerializer = new ObjectWriter() {
    @Override
    public void write(JSONWriter jsonWriter, Object object, Object fieldName, Type fieldType, long features) {
    DescriptorsJSONResult value = (DescriptorsJSONResult) object;
    Objects.requireNonNull(value, "callable of " + fieldName + " can not be null");
    jsonWriter.writeRaw(value.toJSONString());
    }
    };
    com.alibaba.fastjson2.JSON.register(DescriptorsJSONResult.class, descSerializer);
  5. 通过ServerPortExport.json配置描述文件,设置属性serverPort的默认值
    ServerPortExport.json
      {
    "serverPort": {
    "help": "SpringBoot配置,HTTP端口号,默认7700,不建议更改",
    "dftVal": "com.qlangtech.tis.plugin.datax.powerjob.ServerPortExport.dftExportPort():uncache_true"
    }
    }
    ServerPortExport.java
      public static Integer dftExportPort() {
    return ((DefaultExportPortProvider)
    DescriptorsJSONResult.getRootDescInstance()).get();
    }

总结

通过以上步骤,就可以将ServerPortExport根据所在聚合类不同将属性serverPort初始化成不同的默认值。 以此作为一个例子,可以在TIS中相同需求可以推而广之。