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/:AssertSink、AssertSinkFactory、AssertSinkOptions(选项定义)、AssertSinkWriter(校验执行入口);rule/:AssertFieldRule(字段规则模型与规则类型枚举)、AssertRuleParser(HOCON 配置解析)、AssertTableRule、AssertCatalogTableRule(Catalog 元数据规则);excecutor/:AssertExecutor(单条数据校验执行器,被AssertSinkWriter复用);exception/:AssertConnectorErrorCode、AssertConnectorException(统一异常)。
二、引擎支持与功能特性
支持的引擎
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 中定义的常量一一对应):
| 名称 | 类型 | 必填 | 默认值 |
|---|---|---|---|
| rules | ConfigMap | yes | - |
| rules.field_rules | ConfigList | no | - |
| rules.field_rules.field_name | string|ConfigMap | yes | - |
| rules.field_rules.field_type | string | no | - |
| rules.field_rules.field_value | ConfigList | no | - |
| rules.field_rules.field_value.rule_type | string | no | - |
| rules.field_rules.field_value.rule_value | numeric | no | - |
| rules.field_rules.field_value.equals_to | boolean|numeric|string|ConfigList|ConfigMap | no | - |
| rules.row_rules | ConfigList | no | - |
| rules.row_rules.rule_type | string | no | - |
| rules.row_rules.rule_value | string | no | - |
| rules.catalog_table_rule | ConfigMap | no | - |
| rules.catalog_table_rule.primary_key_rule | ConfigMap | no | - |
| rules.catalog_table_rule.primary_key_rule.primary_key_name | string | no | - |
| rules.catalog_table_rule.primary_key_rule.primary_key_columns | ConfigList | no | - |
| rules.catalog_table_rule.constraint_key_rule | ConfigList | no | - |
| rules.catalog_table_rule.constraint_key_rule.constraint_key_name | string | no | - |
| rules.catalog_table_rule.constraint_key_rule.constraint_key_type | string | no | - |
| rules.catalog_table_rule.constraint_key_rule.constraint_key_columns | ConfigList | no | - |
| rules.catalog_table_rule.constraint_key_rule.constraint_key_columns.constraint_key_column_name | string | no | - |
| rules.catalog_table_rule.constraint_key_rule.constraint_key_columns.constraint_key_sort_type | string | no | - |
| rules.catalog_table_rule.column_rule | ConfigList | no | - |
| rules.catalog_table_rule.column_rule.name | string | no | - |
| rules.catalog_table_rule.column_rule.type | string | no | - |
| rules.catalog_table_rule.column_rule.column_length | int | no | - |
| rules.catalog_table_rule.column_rule.nullable | boolean | no | - |
| rules.catalog_table_rule.column_rule.default_value | string | no | - |
| rules.catalog_table_rule.column_rule.comment | comment | no | - |
| rules.table-names | ConfigList | no | - |
| rules.tables_configs | ConfigList | no | - |
| rules.tables_configs.table_path | String | no | - |
| multi_table_sink_replica | int | no | - |
| common-options | no | - |
重要约束:虽然只有
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 文档 中"如何声明受支持的类型"一节的约定。例如
string、int、decimal(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_rules的field_value中。注意源码中rule_value的类型是Double(AssertRule.AssertRule中的private Double ruleValue),因此rule_value需书写为数值。
4.4 rule_value(numeric)
与规则类型配套的值。当rule_type为MIN、MAX、MIN_LENGTH、MAX_LENGTH、MIN_ROW或MAX_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 递归比较、BYTES走Arrays.equals、其余标量类型走value.equals(confValue)。这意味着equals_to对array<int>、map<time, string>、嵌套 row 等复杂值都能做结构化深度相等校验。
4.6 catalog_table_rule(ConfigMap)—— Catalog 表元数据断言
用于断言实际 Catalog 表元数据与用户定义的表元数据完全一致。它由四个子规则组成(对应 AssertCatalogTableRule.java 中的四个字段):
- primary_key_rule:校验主键,包含
primary_key_name(主键名)与primary_key_columns(主键列集合)。源码中的比较逻辑为:若期望主键名为空则跳过名比较;若期望列集合非空,则用CollectionUtils.isEqualCollection做无序集合相等比较; - constraint_key_rule:校验约束键列表,每项包含
constraint_key_name、constraint_key_type(如UNIQUE_KEY)、constraint_key_columns(每列含constraint_key_column_name与constraint_key_sort_type,如ASC)。源码同样使用isEqualCollection与实际的约束键集合比较; - column_rule:按顺序逐列比较,每列声明
name、type、column_length、nullable、default_value、comment。源码isColumnEqual的比较维度包括:列名、数据类型、列长度、scale、是否可空、默认值、注释、源类型,任一维度不一致即抛错,且要求列数完全一致; - table_identifier_rule(源码中位于
AssertConfig.TableIdentifierRule,含catalog_name与table):校验TableIdentifier与期望完全相等。
4.7 table-names(ConfigList)—— 表名存在性断言
用于断言输入数据中确实存在所列的表名。从 AssertSinkWriter.close() 的实现看:Writer 在write阶段会把每行数据的tableId收集进静态集合TABLE_NAMES,close()时若配置了table-names,会用new HashSet<>(期望表名).equals(TABLE_NAMES)做集合相等校验——注意这要求实际收集到的表名集合与期望集合完全相同,不仅仅是"包含关系"。
4.8 tables_configs(ConfigList)—— 多表独立规则
用于为多张输入表定义不同的断言规则,每个元素必须包含table_path。table_path的取值必须与上游 Source 携带的表路径一致(在 Zeta 引擎中,上游表路径一般来自 Source 的tables_configs或schema.table声明)。在多表作业中,每条规则可以独立配置row_rules与field_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 同时支持BATCH与STREAMING两种作业模式,但两条执行路径的时机完全不同,这一点在原文档 sink/Assert.md 的 "Streaming Validation" 一节有明确说明,且与源码实现完全吻合:
- 字段级规则逐行生效:
NOT_NULL、MIN_LENGTH、MAX_LENGTH、MIN、MAX、equals_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):
- Decimal 校验同时检查精度与小数位:
checkDecimalType会校验实际值的scale()必须等于声明的 scale,且precision()必须小于等于声明的 precision,否则类型检查不通过; - 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)是整条校验链路的入口,核心逻辑可以归纳为三步:
- 记录表名与累计行数:把当前行的
tableId加入静态集合TABLE_NAMES;用LONG_ACCUMULATOR.computeIfAbsent(tableName, k -> new LongAccumulator(Long::sum, 0)).accumulate(1)对该表的行数累加 1; - 定位字段规则:单表场景直接取唯一一份
field_rules;多表场景按行的tableId(或catalogTableName)从Map<String, List<AssertFieldRule>>中取出对应规则; - 执行字段规则:调用静态单例
ASSERT_EXECUTOR.fail(element, seaTunnelRowType, rules),一旦Optional<AssertFieldRule>非空,立即抛出AssertConnectorException(RULE_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/MAX用doubleValue()比较、MIN_LENGTH/MAX_LENGTH用字符串长度比较);随后若配置了equals_to且值非空,再执行compareValue做结构化深度比较。
7.3 行数规则:close 时的一次性断言
AssertSinkWriter.close()中,row_rules(MIN_ROW/MAX_ROW)只在该 Writer 负责的表上执行一次:从LONG_ACCUMULATOR取出该表累计行数,若count > MAX_ROW或count < 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_configs与table_list; - 2.3.10:新增 Assert options;重构 connector common options;
- 2.3.12:将元数据 schema 引入 Catalog 表(add metadata schema into catalog table)。
九、常见误区与使用建议
- 行数规则不是逐行触发的:
MIN_ROW/MAX_ROW只在 Writer 关闭时执行一次,且统计的是该 Writer 自创建以来的累积行数。若在流式作业中期望"每个 checkpoint 窗口校验行数",当前版本无法直接满足,需要源码级改动; equals_to对 null 字段无效:null 类型字段请用rule_type = NULL/NOT_NULL控制;table-names是集合相等而非包含:期望表名集合必须与实际接收到的表名集合完全一致,否则 close 时失败;catalog_table_rule是全维度逐列比对:列数、列名、类型、长度、scale、可空性、默认值、注释、源类型任一不一致都会使作业失败,配置前请确认上游 Catalog 元数据;- 至少配置一条有意义的规则:只有
rules是必填,但空规则集意味着"无事可校验",不会起到任何质量保障作用; 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),仅供参考