SeaTunnel Assert Sink 连接器完全指南:用规则化校验构建可信赖的数据管道
2026/9/16 11:18:11 网站建设 项目流程

SeaTunnel Assert Sink 连接器完全指南:用规则化校验构建可信赖的数据管道

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

本指南围绕 Apache SeaTunnel 的 Assert Sink 连接器展开。Assert 是一个"零外部依赖"的终端校验型 Sink,它不向任何外部系统写数据,而是按照用户在配置中声明的规则对管道输出进行断言:行数是否落在期望区间、字段类型与取值是否合法、字段是否可空、多表场景下各表名与元数据是否与预期一致。读完本文,你将掌握 Assert 全部规则体系(field_rules/row_rules/catalog_table_rule/tables_configs)的配置方法与流式/批式语义,并能结合源码理解每条规则的执行时机与失败行为,从而把数据质量校验直接内嵌到 SeaTunnel 作业中。

一、Assert 是什么:一个没有外部系统的终端 Sink

Assert 是 SeaTunnel 连接器体系中的校验型 Sink 连接器,其定位与 Kafka、JDBC 等写入型 Sink 完全不同:它没有外部系统可写,唯一的作用是"接住"上游数据并对它们做校验。官方文档的描述是:它通过用户自定义规则检查行数(row count)、字段类型(field type)、字段值(field value)以及 Catalog 表元数据(catalog table metadata),一旦实际数据与规则不匹配,作业就会失败(job fails)。

这个定位让 Assert 在以下场景中极具价值:

  • 管道自检(Pipeline Self-check):在开发或回归测试阶段,把 Assert 挂在管道末端,验证 FakeSource 或任意上游产出的数据是否符合预期,无需准备下游数据库;
  • 中间结果验证:配合 Transform 使用,校验转换后的中间结果,避免"脏数据流向下游";
  • 多表作业校验:对一个作业中的多张输入表分别断言行数与字段规则;
  • 流式作业的累积校验:在流式模式下对 Sink Writer 生命周期内累计接收的行数做最终校验。

从源码结构看,Assert 连接器由以下几部分构成(seatunnel-connectors-v2/connector-assert):

  • sink/AssertSinkAssertSinkFactoryAssertSinkOptions(选项定义)、AssertSinkWriter(校验执行入口);
  • rule/AssertFieldRule(字段规则模型与规则类型枚举)、AssertRuleParser(HOCON 配置解析)、AssertTableRuleAssertCatalogTableRule(Catalog 元数据规则);
  • excecutor/AssertExecutor(单条数据校验执行器,被AssertSinkWriter复用);
  • exception/AssertConnectorErrorCodeAssertConnectorException(统一异常)。

二、引擎支持与功能特性

支持的引擎

Assert 连接器同时支持三大执行引擎:

Spark / Flink / Seatunnel Zeta

功能特性一览

特性支持
exactly-once
cdc
batch
stream
多表写入(multiple table write)
timer flush

需要说明的是:Assert 是终端校验型 Sink,本身不涉及"恰好一次"之类的语义保证;它也不将UPDATE/DELETE行类型解释为 CDC 操作,而是对接收到的每一行都执行规则断言(见原文档 sink/Assert.md 中的 tip 说明)。

三、选项总览:整棵规则树

Assert 连接器只有顶层一个必填参数rules,其余全部为嵌套可选规则。完整选项表如下(字段名与 AssertConfig.java 中定义的常量一一对应):

名称类型必填默认值
rulesConfigMapyes-
rules.field_rulesConfigListno-
rules.field_rules.field_namestring|ConfigMapyes-
rules.field_rules.field_typestringno-
rules.field_rules.field_valueConfigListno-
rules.field_rules.field_value.rule_typestringno-
rules.field_rules.field_value.rule_valuenumericno-
rules.field_rules.field_value.equals_toboolean|numeric|string|ConfigList|ConfigMapno-
rules.row_rulesConfigListno-
rules.row_rules.rule_typestringno-
rules.row_rules.rule_valuestringno-
rules.catalog_table_ruleConfigMapno-
rules.catalog_table_rule.primary_key_ruleConfigMapno-
rules.catalog_table_rule.primary_key_rule.primary_key_namestringno-
rules.catalog_table_rule.primary_key_rule.primary_key_columnsConfigListno-
rules.catalog_table_rule.constraint_key_ruleConfigListno-
rules.catalog_table_rule.constraint_key_rule.constraint_key_namestringno-
rules.catalog_table_rule.constraint_key_rule.constraint_key_typestringno-
rules.catalog_table_rule.constraint_key_rule.constraint_key_columnsConfigListno-
rules.catalog_table_rule.constraint_key_rule.constraint_key_columns.constraint_key_column_namestringno-
rules.catalog_table_rule.constraint_key_rule.constraint_key_columns.constraint_key_sort_typestringno-
rules.catalog_table_rule.column_ruleConfigListno-
rules.catalog_table_rule.column_rule.namestringno-
rules.catalog_table_rule.column_rule.typestringno-
rules.catalog_table_rule.column_rule.column_lengthintno-
rules.catalog_table_rule.column_rule.nullablebooleanno-
rules.catalog_table_rule.column_rule.default_valuestringno-
rules.catalog_table_rule.column_rule.commentcommentno-
rules.table-namesConfigListno-
rules.tables_configsConfigListno-
rules.tables_configs.table_pathStringno-
multi_table_sink_replicaintno-
common-optionsno-

重要约束:虽然只有rules是必填项,但嵌套的规则块是可选的,至少要配置一条有意义的规则,否则 Sink 将无内容可校验。

从源码看,顶层rules选项的定义位于 AssertSinkOptions.java,其类型为Map<String, Object>,无默认值、必填;而AssertSinkOptions本身继承自SinkConnectorCommonOptions,因此multi_table_sink_replica与 common options 均来自公共 Sink 选项体系。

四、核心规则详解

4.1 rules(ConfigMap,必填)

规则的唯一必填顶层块,描述期望数据。每个规则代表一种校验:字段校验、行数校验、表名校验或 Catalog 表校验。

4.2 field_rules(ConfigList)—— 字段级校验

当需要校验字段类型、空值约束、值域范围、字符串长度或精确取值时使用。每条 field rule 由三部分构成:

  • field_name(string,必填):被校验的字段名。从 AssertExecutor.java 的实现看,执行时会先在SeaTunnelRowType中按字段名查找索引,如果字段不存在会直接抛出IllegalArgumentException("Field name %s not found in row type %s");
  • field_type(string | ConfigMap,可选):字段类型声明,声明方式需遵循 schema-feature 文档 中"如何声明受支持的类型"一节的约定。例如stringintdecimal(30, 8)array<int>map<time, string>,嵌套行类型则直接以 ConfigMap 形式声明;
  • field_value(ConfigList,可选):字段值的校验规则列表,可叠加多条规则,例如同时校验"非空"与"长度上下限"。

4.3 rule_type(string)—— 支持的规则类型全集

rule_type是值校验的核心。目前支持的枚举定义在 AssertFieldRule.java 的AssertRuleType中,与官方文档完全一致:

rule_type含义是否需要 rule_value
NOT_NULL值不能为 null
NULL值可以为 null
MIN数据的最小值(数值下界,含等号)
MAX数据的最大值(数值上界,含等号)
MIN_LENGTH字符串数据的最小长度
MAX_LENGTH字符串数据的最大长度
MIN_ROW最小行数
MAX_ROW最大行数

其中MIN_ROW/MAX_ROW属于行级规则,应放在row_rules块中;其余属于字段值规则,放在field_rulesfield_value中。注意源码中rule_value的类型是DoubleAssertRule.AssertRule中的private Double ruleValue),因此rule_value需书写为数值。

4.4 rule_value(numeric)

与规则类型配套的值。当rule_typeMINMAXMIN_LENGTHMAX_LENGTHMIN_ROWMAX_ROW时,必须为rule_value赋值。

4.5 equals_to(boolean | numeric | string | ConfigList | ConfigMap)—— 精确值比较

equals_to用于比较字段实际值是否等于配置的期望值,支持所有 SeaTunnel 类型(类型清单见 schema-feature 文档 的"当前支持哪些类型"一节)。对于复杂类型,需使用与上游数据一致的 HOCON 值形态。文档给出的示例:某字段是包含三个字段的 row 类型,声明为{a = array<string>, b = map<string, decimal(30, 2)>, c={c_0 = int, b = string}},则期望值可以写成[["a", "b"], { k0 = 9999.99, k1 = 111.11 }, [123, "abcd"]]

两个重要注意点:

  • 定义字段值的方式与 FakeSource 连接器 中"定制数据内容"的方式一致;
  • equals_to不能应用于null类型字段,此时请改用规则类型NULL进行校验,例如{rule_type = NULL}

从实现看,equals_to的比对并非简单的字符串相等:AssertExecutor.compareValue会先把配置值通过JsonToRowConverters转换为 SeaTunnel 类型对象,再按类型分发比对——ROW走逐字段递归比较、ARRAY走逐元素比较(长度也必须一致)、MAP比较 key 是否都存在且逐 value 递归比较、BYTESArrays.equals、其余标量类型走value.equals(confValue)。这意味着equals_toarray<int>map<time, string>、嵌套 row 等复杂值都能做结构化深度相等校验。

4.6 catalog_table_rule(ConfigMap)—— Catalog 表元数据断言

用于断言实际 Catalog 表元数据与用户定义的表元数据完全一致。它由四个子规则组成(对应 AssertCatalogTableRule.java 中的四个字段):

  1. primary_key_rule:校验主键,包含primary_key_name(主键名)与primary_key_columns(主键列集合)。源码中的比较逻辑为:若期望主键名为空则跳过名比较;若期望列集合非空,则用CollectionUtils.isEqualCollection无序集合相等比较;
  2. constraint_key_rule:校验约束键列表,每项包含constraint_key_nameconstraint_key_type(如UNIQUE_KEY)、constraint_key_columns(每列含constraint_key_column_nameconstraint_key_sort_type,如ASC)。源码同样使用isEqualCollection与实际的约束键集合比较;
  3. column_rule:按顺序逐列比较,每列声明nametypecolumn_lengthnullabledefault_valuecomment。源码isColumnEqual的比较维度包括:列名、数据类型、列长度、scale、是否可空、默认值、注释、源类型,任一维度不一致即抛错,且要求列数完全一致;
  4. table_identifier_rule(源码中位于AssertConfig.TableIdentifierRule,含catalog_nametable):校验TableIdentifier与期望完全相等。

4.7 table-names(ConfigList)—— 表名存在性断言

用于断言输入数据中确实存在所列的表名。从 AssertSinkWriter.close() 的实现看:Writer 在write阶段会把每行数据的tableId收集进静态集合TABLE_NAMESclose()时若配置了table-names,会用new HashSet<>(期望表名).equals(TABLE_NAMES)集合相等校验——注意这要求实际收集到的表名集合与期望集合完全相同,不仅仅是"包含关系"。

4.8 tables_configs(ConfigList)—— 多表独立规则

用于为多张输入表定义不同的断言规则,每个元素必须包含table_pathtable_path的取值必须与上游 Source 携带的表路径一致(在 Zeta 引擎中,上游表路径一般来自 Source 的tables_configsschema.table声明)。在多表作业中,每条规则可以独立配置row_rulesfield_rules

4.9 multi_table_sink_replica(int)

多表 Sink 的公共选项——每个表所需的 Sink 副本(replica)数量。仅在作业需要"每张表多个 Sink 副本"时配置。其余 Sink 公共参数参见 Sink 公共选项文档。

五、规则匹配语义与执行时机

5.1 三类规则的作用对象

  • row_rules:检查 Assert Sink 接收到的行数
  • field_rules:检查每一行中指定字段的值;
  • tables_configs:用于多表作业,table_path必须与上游 Source 携带的表路径匹配;
  • equals_to:将实际字段值与配置的期望值比较,复杂值(array、map、row)需使用与源数据相同的 HOCON 值形态。

5.2 批式与流式下的执行时机(重要)

Assert 同时支持BATCHSTREAMING两种作业模式,但两条执行路径的时机完全不同,这一点在原文档 sink/Assert.md 的 "Streaming Validation" 一节有明确说明,且与源码实现完全吻合:

  • 字段级规则逐行生效NOT_NULLMIN_LENGTHMAX_LENGTHMINMAXequals_to等在每行数据到达 Sink Writer 时立即检查(AssertSinkWriter.write中调用ASSERT_EXECUTOR.fail(...),命中失败规则即抛出AssertConnectorException,错误码为RULE_VALIDATION_FAILED,消息形如row :<row> fail rule: <rule>);
  • 行数规则仅在 Writer 关闭时执行一次MIN_ROW/MAX_ROW在 Sink Writerclose()时(作业结束、savepoint 或失败时)恰好执行一次,针对的是该 Writer 实例自创建以来累计的行数——而非每个 checkpoint 窗口的行数,也不会在 checkpoint 之间重置。累计计数通过静态ConcurrentHashMap<String, LongAccumulator>LONG_ACCUMULATOR,以表名为 key)实现,天然支持跨 checkpoint 累积。

如果你的需求是"按 checkpoint 窗口做行数校验",这超出了当前文档与实现的能力范围,需要修改源码才能支持。

5.3 多表场景下的 close 语义

AssertSinkWriter.close()对行数规则的校验逻辑有一个细节:在多表作业中,只有当assertRowRules.size() == 1或规则的 key 与当前 Writer 的catalogTableName相等时才会执行该校验——也就是说,每个 Writer 只对自己负责的那张表做行数断言。这一点由测试 AssertSinkWriterCloseTest.java 中的两个用例直接验证:

  • testCloseOnlyAssertsOwnTableWhenMultipleTables:Writer A 收到 1 行、表 B 从未收到行,A 的MIN_ROW=1校验必须通过,不能被其他 Writer 的进度影响;
  • testCloseThrowsWhenOwnTableRuleNotMet:A 自己只写了 1 行但MIN_ROW=2,close 时必须抛AssertConnectorException

六、实战示例

6.1 简单示例:行数 + 字段规则 + Catalog 元数据

下面的作业验证管道输出行数在 5~10 之间,且name字段非空、字符串长度在 5~10,age字段非空、精确等于 23、值域在 32767~2147483647,同时断言 Catalog 表的主键、唯一约束与列元数据。

Assert { rules = { row_rules = [ { rule_type = MAX_ROW rule_value = 10 }, { rule_type = MIN_ROW rule_value = 5 } ], field_rules = [{ field_name = name field_type = string field_value = [ { rule_type = NOT_NULL }, { rule_type = MIN_LENGTH rule_value = 5 }, { rule_type = MAX_LENGTH rule_value = 10 } ] }, { field_name = age field_type = int field_value = [ { rule_type = NOT_NULL equals_to = 23 }, { rule_type = MIN rule_value = 32767 }, { rule_type = MAX rule_value = 2147483647 } ] } ] catalog_table_rule { primary_key_rule = { primary_key_name = "primary key" primary_key_columns = ["id"] } constraint_key_rule = [ { constraint_key_name = "unique_name" constraint_key_type = UNIQUE_KEY constraint_key_columns = [ { constraint_key_column_name = "id" constraint_key_sort_type = ASC } ] } ] column_rule = [ { name = "id" type = bigint }, { name = "name" type = string }, { name = "age" type = int } ] } } }

注意:catalog_table_rule是对Catalog 元数据的断言,示例中的column_rule必须与实际管道中表的列定义逐维一致(名称、类型、长度、可空性、默认值、注释等),否则作业会以CATALOG_TABLE_FAILED错误码失败。

6.2 复杂示例:equals_to 全类型精确比对

下面的完整作业先用 FakeSource 产出一条包含 null、string、boolean、整数族、浮点族、decimal、date/timestamp/time、bytes、array、map(含嵌套 map)与嵌套 row 的记录,再用 Assert 的equals_to逐字段做精确比对。这是equals_to最典型的全类型验证场景。FakeSource 的用法参见 FakeSource 文档。

source { FakeSource { row.num = 1 schema = { fields { c_null = "null" c_string = string c_boolean = boolean c_tinyint = tinyint c_smallint = smallint c_int = int c_bigint = bigint c_float = float c_double = double c_decimal = "decimal(30, 8)" c_date = date c_timestamp = timestamp c_time = time c_bytes = bytes c_array = "array<int>" c_map = "map<time, string>" c_map_nest = "map<string, {c_int = int, c_string = string}>" c_row = { c_null = "null" c_string = string c_boolean = boolean c_tinyint = tinyint c_smallint = smallint c_int = int c_bigint = bigint c_float = float c_double = double c_decimal = "decimal(30, 8)" c_date = date c_timestamp = timestamp c_time = time c_bytes = bytes c_array = "array<int>" c_map = "map<string, string>" } } } rows = [ { kind = INSERT fields = [ null, "AAA", false, 1, 1, 333, 323232, 3.1, 9.33333, 99999.99999999, "2012-12-21", "2012-12-21T12:34:56", "12:34:56", "bWlJWmo=", [0, 1, 2], "{ 12:01:26 = v0 }", { k1 = [123, "BBB-BB"]}, [ null, "AAA", false, 1, 1, 333, 323232, 3.1, 9.33333, 99999.99999999, "2012-12-21", "2012-12-21T12:34:56", "12:34:56", "bWlJWmo=", [0, 1, 2], { k0 = v0 } ] ] } ] plugin_output = "fake" } } sink{ Assert { plugin_input = "fake" rules = { row_rules = [ { rule_type = MAX_ROW rule_value = 1 }, { rule_type = MIN_ROW rule_value = 1 } ], field_rules = [ { field_name = c_null field_type = "null" field_value = [ { rule_type = NULL } ] }, { field_name = c_string field_type = string field_value = [ { rule_type = NOT_NULL equals_to = "AAA" } ] }, { field_name = c_boolean field_type = boolean field_value = [ { rule_type = NOT_NULL equals_to = false } ] }, { field_name = c_tinyint field_type = tinyint field_value = [ { rule_type = NOT_NULL equals_to = 1 } ] }, { field_name = c_smallint field_type = smallint field_value = [ { rule_type = NOT_NULL equals_to = 1 } ] }, { field_name = c_int field_type = int field_value = [ { rule_type = NOT_NULL equals_to = 333 } ] }, { field_name = c_bigint field_type = bigint field_value = [ { rule_type = NOT_NULL equals_to = 323232 } ] }, { field_name = c_float field_type = float field_value = [ { rule_type = NOT_NULL equals_to = 3.1 } ] }, { field_name = c_double field_type = double field_value = [ { rule_type = NOT_NULL equals_to = 9.33333 } ] }, { field_name = c_decimal field_type = "decimal(30, 8)" field_value = [ { rule_type = NOT_NULL equals_to = 99999.99999999 } ] }, { field_name = c_date field_type = date field_value = [ { rule_type = NOT_NULL equals_to = "2012-12-21" } ] }, { field_name = c_timestamp field_type = timestamp field_value = [ { rule_type = NOT_NULL equals_to = "2012-12-21T12:34:56" } ] }, { field_name = c_time field_type = time field_value = [ { rule_type = NOT_NULL equals_to = "12:34:56" } ] }, { field_name = c_bytes field_type = bytes field_value = [ { rule_type = NOT_NULL equals_to = "bWlJWmo=" } ] }, { field_name = c_array field_type = "array<int>" field_value = [ { rule_type = NOT_NULL equals_to = [0, 1, 2] } ] }, { field_name = c_map field_type = "map<time, string>" field_value = [ { rule_type = NOT_NULL equals_to = "{ 12:01:26 = v0 }" } ] }, { field_name = c_map_nest field_type = "map<string, {c_int = int, c_string = string}>" field_value = [ { rule_type = NOT_NULL equals_to = { k1 = [123, "BBB-BB"] } } ] }, { field_name = c_row field_type = { c_null = "null" c_string = string c_boolean = boolean c_tinyint = tinyint c_smallint = smallint c_int = int c_bigint = bigint c_float = float c_double = double c_decimal = "decimal(30, 8)" c_date = date c_timestamp = timestamp c_time = time c_bytes = bytes c_array = "array<int>" c_map = "map<string, string>" } field_value = [ { rule_type = NOT_NULL equals_to = [ null, "AAA", false, 1, 1, 333, 323232, 3.1, 9.33333, 99999.99999999, "2012-12-21", "2012-12-21T12:34:56", "12:34:56", "bWlJWmo=", [0, 1, 2], { k0 = v0 } ] } ] } ] } } }

从实现角度补充两个可验证的细节(见 AssertExecutor.java):

  1. Decimal 校验同时检查精度与小数位checkDecimalType会校验实际值的scale()必须等于声明的 scale,且precision()必须小于等于声明的 precision,否则类型检查不通过;
  2. NULL 类型字段的强约束checkType中若字段类型为SqlType.NULL而值非空,直接返回 false;反之若字段声明了具体类型而值为 null,类型检查本身放行,是否允许 null 交由rule_type = NULL/NOT_NULL决定——这就是为什么c_null字段要用rule_type = NULL校验,而不能用equals_to

6.3 多表断言:每张表独立的行数与字段规则

下面的作业在同一个 Job 中处理两张表:test.table1(16 行,c_int/c_bigint)与test.table2(17 行,c_string/c_tinyint),Assert 通过tables_configs对每张表分别断言行数与字段非空。

env { parallelism = 1 job.mode = BATCH } source { FakeSource { tables_configs = [ { row.num = 16 schema { table = "test.table1" fields { c_int = int c_bigint = bigint } } }, { row.num = 17 schema { table = "test.table2" fields { c_string = string c_tinyint = tinyint } } } ] } } transform { } sink { Assert { rules = { tables_configs = [ { table_path = "test.table1" row_rules = [ { rule_type = MAX_ROW rule_value = 16 }, { rule_type = MIN_ROW rule_value = 16 } ], field_rules = [{ field_name = c_int field_type = int field_value = [ { rule_type = NOT_NULL } ] }, { field_name = c_bigint field_type = bigint field_value = [ { rule_type = NOT_NULL } ] }] }, { table_path = "test.table2" row_rules = [ { rule_type = MAX_ROW rule_value = 17 }, { rule_type = MIN_ROW rule_value = 17 } ], field_rules = [{ field_name = c_string field_type = string field_value = [ { rule_type = NOT_NULL } ] }, { field_name = c_tinyint field_type = tinyint field_value = [ { rule_type = NOT_NULL } ] }] } ] } } }

6.4 流式校验示例:累积行数窗口 + 字段值域

下面的流式作业以 60 秒为 checkpoint 间隔运行,FakeSource 产出 1000 行数据。Assert 在 Writer 关闭时对累积行数50 ≤ total rows ≤ 5000的校验(注意不是按 checkpoint 窗口),同时对每行的age字段执行非空与0 ≤ age ≤ 150的逐行校验。

env { parallelism = 1 job.mode = "STREAMING" checkpoint.interval = 60000 } source { FakeSource { row.num = 1000 schema { fields { name = string age = int } } plugin_output = "stream_data" } } sink { Assert { plugin_input = "stream_data" rules = { row_rules = [ { rule_type = MIN_ROW rule_value = 50 }, { rule_type = MAX_ROW rule_value = 5000 } ], field_rules = [{ field_name = age field_type = int field_value = [ { rule_type = NOT_NULL }, { rule_type = MIN rule_value = 0 }, { rule_type = MAX rule_value = 150 } ] }] } } }

七、源码级原理剖析

7.1 单行校验链路:AssertSinkWriter → AssertExecutor

AssertSinkWriter.write(SeaTunnelRow)(AssertSinkWriter.java)是整条校验链路的入口,核心逻辑可以归纳为三步:

  1. 记录表名与累计行数:把当前行的tableId加入静态集合TABLE_NAMES;用LONG_ACCUMULATOR.computeIfAbsent(tableName, k -> new LongAccumulator(Long::sum, 0)).accumulate(1)对该表的行数累加 1;
  2. 定位字段规则:单表场景直接取唯一一份field_rules;多表场景按行的tableId(或catalogTableName)从Map<String, List<AssertFieldRule>>中取出对应规则;
  3. 执行字段规则:调用静态单例ASSERT_EXECUTOR.fail(element, seaTunnelRowType, rules),一旦Optional<AssertFieldRule>非空,立即抛出AssertConnectorExceptionRULE_VALIDATION_FAILED)。

7.2 规则执行器:AssertExecutor 的判定顺序

AssertExecutor(AssertExecutor.java)对每个字段的判定顺序是:先类型后值

  • checkType:按 SQL 类型分发——ROW递归检查每个子字段、ARRAY检查每个元素、MAP同时检查 key 与 value、DECIMAL检查精度与小数位、向量类型要求ByteBuffer、其余类型要求value.getClass().equals(fieldType.getTypeClass())
  • checkValue:逐条执行field_value中的规则。rule_type不为空时先按checkAssertRule的 switch 判定(NULL/NOT_NULL判空、MIN/MAXdoubleValue()比较、MIN_LENGTH/MAX_LENGTH用字符串长度比较);随后若配置了equals_to且值非空,再执行compareValue做结构化深度比较。

7.3 行数规则:close 时的一次性断言

AssertSinkWriter.close()中,row_rulesMIN_ROW/MAX_ROW)只在该 Writer 负责的表上执行一次:从LONG_ACCUMULATOR取出该表累计行数,若count > MAX_ROWcount < MIN_ROW,抛出携带实际行数的AssertConnectorException(消息形如row num :<count> fail rule: <rule>)。table-names的集合相等校验也在此处完成。

7.4 Catalog 元数据断言:AssertCatalogTableRule

AssertCatalogTableRule.checkRule(CatalogTable)(AssertCatalogTableRule.java)依次校验主键、约束键、列与表标识符,任何不一致都以CATALOG_TABLE_FAILED错误码失败。列比较isColumnEqual覆盖了名称、数据类型、列长度、scale、可空性、默认值、注释、源类型共 8 个维度。

八、版本演进时间线(Changelog)

Assert 连接器自 2.2.0-beta 引入以来持续演进,各版本的关键能力如下(依据 connector-assert changelog 整理):

  • 2.2.0-beta:将 Assert Sink 加入 API 草案(add assert sink to Api draft),随后以统一工厂(Source/Sink Factory)形式完善连接器框架;
  • 2.3.0 / 2.3.0-beta:为 Assert 增加 Sink 工厂并统一异常体系(Unified exception),保证factoryIdentifier与插件名一致;
  • 2.3.1:配合 SimpleSQL 等 Transform 演进,完善多模块构建;
  • 2.3.4(能力密集版本):支持全数据类型的字段类型断言与字段值相等断言(field type assert & field value equality assert for full data types);支持检查 Decimal 类型的 precision 与 scale;新增table-names以支持 FakeSource/Assert 的多表产出与断言;schema 支持配置 column/primaryKey/constraintKey;
  • 2.3.5:修复 DateTime 工具相关问题;
  • 2.3.6:支持在 Sink 选项中使用上游表占位符并自动替换;
  • 2.3.7:增加多表 Sink 选项检查;
  • 2.3.8支持多表校验(Assert support multi-table check);
  • 2.3.9:支持带时区偏移的时间戳(timestamp with timezone offset);优化 Assert Sink 校验方法;统一tables_configstable_list
  • 2.3.10:新增 Assert options;重构 connector common options;
  • 2.3.12:将元数据 schema 引入 Catalog 表(add metadata schema into catalog table)。

九、常见误区与使用建议

  1. 行数规则不是逐行触发的MIN_ROW/MAX_ROW只在 Writer 关闭时执行一次,且统计的是该 Writer 自创建以来的累积行数。若在流式作业中期望"每个 checkpoint 窗口校验行数",当前版本无法直接满足,需要源码级改动;
  2. equals_to对 null 字段无效:null 类型字段请用rule_type = NULL/NOT_NULL控制;
  3. table-names是集合相等而非包含:期望表名集合必须与实际接收到的表名集合完全一致,否则 close 时失败;
  4. catalog_table_rule是全维度逐列比对:列数、列名、类型、长度、scale、可空性、默认值、注释、源类型任一不一致都会使作业失败,配置前请确认上游 Catalog 元数据;
  5. 至少配置一条有意义的规则:只有rules是必填,但空规则集意味着"无事可校验",不会起到任何质量保障作用;
  6. table_path必须与上游一致:多表场景下tables_configs中的table_path必须与 Source 携带的表路径精确匹配。

十、相关资源

  • Assert 连接器主文档:docs/en/connectors/sink/Assert.md
  • 变更日志:docs/en/connectors/changelog/connector-assert.md
  • 源码目录:seatunnel-connectors-v2/connector-assert(含sink/rule/excecutor/exception/子包)
  • 单元测试:AssertSinkWriterCloseTest.java(多表 close 语义)、AssertExecutorTest.java
  • 配套的 FakeSource 文档:docs/en/connectors/source/FakeSource.md(equals_to值形态定义方式与其一致)
  • Sink 公共选项:docs/en/connectors/common-options/sink-common-options.md
  • Schema 类型声明指南:docs/en/introduction/concepts/schema-feature.md

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询