
SeaTunnel Kafka 连接器 COMPATIBLE_KAFKA_CONNECT_JSON 格式实战消费 Kafka Connect JDBC 与 Debezium 数据【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel 的 Kafka 连接器除了原生 JSON/TEXT 格式外还内置了COMPATIBLE_KAFKA_CONNECT_JSON反序列化格式专门用于解析通过 Kafka Connect Source 抽取到 Kafka 中的数据尤其是 Kafka Connect JDBC Source 和 Debezium 这类以 Connect JSON Envelope带 schema 包装形式输出的消息。读完本文你将掌握该格式的完整作业配置、两个关键 schema 开关参数的作用以及 SeaTunnel 底层JsonConverter反射调用、payload 提取与事件时间挂载的实现机制。适用场景识别 Kafka Connect 风格的 JSON 消息Kafka 中大量数据并非由业务系统直接写入而是经由 Kafka Connect 生态采集而来。这类消息的典型特征是使用了 Kafka Connect 的 JSON Envelope 结构——消息体中除了业务数据外还带有schema字段描述字段类型Kafka Connect JDBC Source抽取的关系型数据库全量/增量数据Debezium基于 Kafka Connect 实现产出的 CDC 数据流其他以 Connect JSON Converterschemas.enabletrue输出格式的 Source。格式文档明确指出SeaTunnel 的 Kafka 连接器支持解析这类数据。在 SeaTunnel 配置中只需在 Kafka Source 里将format设为COMPATIBLE_KAFKA_CONNECT_JSON即可启用无需自行处理 Envelope 包装。完整作业配置从 Kafka 读取 Connect 数据写入 MySQL下面是 官方文档 中给出的完整可运行配置场景为Kafka Connect JDBC Source 将数据库表抽取到 topicjdbc_source_recordSeaTunnel 以 Batch 模式消费该 topic 并写入 MySQL 的jdbc_sink表env { parallelism 1 job.mode BATCH } source { Kafka { bootstrap.servers localhost:9092 topic jdbc_source_record plugin_output kafka_table start_mode earliest schema { fields { id int name string description string weight string } }, format COMPATIBLE_KAFKA_CONNECT_JSON } } sink { Jdbc { driver com.mysql.cj.jdbc.Driver url jdbc:mysql://localhost:3306/seatunnel user st_user password seatunnel generate_sink_sql true database seatunnel table jdbc_sink primary_keys [id] } }配置要点说明配置项说明bootstrap.serversKafka 集群地址必填topic待消费的 topic本例为 Kafka Connect Source 的产出 topicstart_mode earliest从最早 offset 开始消费配合 Batch 模式可完整回放历史数据schema该格式下必填。声明源数据的字段名与类型SeaTunnel 据此构造CatalogTable和SeaTunnelRowType用于把 JSON 节点逐字段转换为SeaTunnelRowformat COMPATIBLE_KAFKA_CONNECT_JSON启用 Kafka Connect 兼容 JSON 反序列化plugin_output定义作业中的表名供下游引用Sink 侧使用 Jdbc 连接器generate_sink_sql true让 SeaTunnel 依据 schema 自动生成 INSERT SQLprimary_keys [id]保证重复消费时按主键覆盖而非重复插入。两个可配置的 schema 开关参数COMPATIBLE_KAFKA_CONNECT_JSON格式提供两个布尔选项定义在 KafkaConnectJsonFormatOptions参数类型默认值说明key_converter_schema_enabledbooleantrue是否按带 schema envelope方式解析消息 Key。即 Key 是{schema: {...}, payload: ...}还是纯值value_converter_schema_enabledbooleantrue是否按带 schema envelope方式解析消息 Value对应 Kafka Connect 中JsonConverter的schemas.enable设置这两个开关的语义直接对齐 Kafka Connect 官方JsonConverter的schemas.enable配置若上游 Connect Source 使用默认配置schemas.enabletrue消息带 schema 包装保持默认true即可若上游显式配置了schemas.enablefalse或仅 Key 关闭了 envelope则应把对应开关改为falseSeaTunnel 底层会切换为不带 envelope 的解析路径。在 KafkaSourceConfig 的createDeserializationSchema方法中可以确认这两个参数被读取后直接传入反序列化 schema 构造器case COMPATIBLE_KAFKA_CONNECT_JSON: Boolean keySchemaEnable readonlyConfig.get(KafkaConnectJsonFormatOptions.KEY_CONVERTER_SCHEMA_ENABLED); Boolean valueSchemaEnable readonlyConfig.get(KafkaConnectJsonFormatOptions.VALUE_CONVERTER_SCHEMA_ENABLED); schema new CompatibleKafkaConnectDeserializationSchema( catalogTable, keySchemaEnable, valueSchemaEnable, false, false); break;另外需要注意从源码的分支结构看当format不是NATIVE且未配置schema时SeaTunnel 会退化为按占位符分隔的文本解析TextDeserializationSchema。因此使用COMPATIBLE_KAFKA_CONNECT_JSON时schema是实质性必填项缺少它将导致字段无法正确映射。反序列化原理复用 Kafka 官方 JsonConverter核心实现位于 seatunnel-format-compatible-connect-json 模块其 pom 声明依赖 Kafka 3.4.0 的kafka-clients与connect-json均为 provided scope由 Kafka 连接器运行时提供。主类 CompatibleKafkaConnectDeserializationSchema 的deserialize(ConsumerRecord, Collector)方法完整处理流程如下初始化 JsonConvertertryInitConverter双重检查锁保证线程安全分别创建 key/value 两个org.apache.kafka.connect.json.JsonConverter并按keySchemaEnable/valueSchemaEnable开关通过反射绑定不同的转换方法带 envelopeconvertToJsonWithEnvelope(Schema, Object)不带 envelopeconvertToJsonWithoutEnvelope(Schema, Object)。转换为 SinkRecordconvertToSinkRecord先把原始ConsumerRecordbyte[], byte[]的 key/value 经JsonConverter.toConnectData转为 Connect 结构化数据再封装成SinkRecord携带 topic、partition、offset、timestamp 等元数据还原出与 Connect Sink 侧一致的记录形态。提取 payload 并展开调用 envelope 转换方法得到 JSON 节点后取出其中的payload字段若 payload 是数组批量消息逐个元素转换为SeaTunnelRow输出否则按单条记录转换。字段映射与元数据挂载借助seatunnel-format-json模块的JsonToRowConverters将 JSON 节点按用户声明的schema映射为SeaTunnelRowattachEventTime将 Kafka 消息时间戳作为event_time元数据写入行选项仅在未显式设置时供下游 watermark / 事件时间处理使用若作业绑定了CatalogTable还会把表路径写入行的tableId支撑多表路由所有行统一标记为RowKind.INSERT。异常处理单个 JSON 节点转换失败时抛出带原始报文上下文的CommonError.jsonOperationError外层 KafkaRecordEmitter 依据 Kafka 连接器的错误处理策略决定跳过该消息打 WARN 日志还是让作业失败。同包下的 NativeKafkaConnectDeserializationSchema 服务于NATIVE格式它不做 Envelope 解析而是把 partition、offset、key、value、timestamp、headers 等原始字段打包成一个 Map 输出适合需要保留完整 Kafka 元信息的场景。Kafka 消息头Header字段透传COMPATIBLE_KAFKA_CONNECT_JSON还支持将 Kafka 消息头作为额外字段注入输出行。从 KafkaRecordEmitter.emitRecord 的实现看当元数据中配置了kafkaHeaderFields时emitter 会包装一个头注入 Collector在每条记录收集时把指定 header 的 UTF-8 值追加到行字段末尾header 不存在则填 null。KafkaSourceConfig 中也有对应逻辑仅当format COMPATIBLE_KAFKA_CONNECT_JSON且存在 header 字段时才会记录这些字段。相关行为可由测试用例 KafkaSourceConfigTest#testKafkaHeaderFieldsExtendsSchemaForCompatibleKafkaConnectJsonFormat 印证——header 字段会自动扩展进输出 schema。限制与注意事项Sink 侧限制从 KafkaSinkWriter 的校验逻辑看COMPATIBLE_KAFKA_CONNECT_JSON以及NATIVE、COMPATIBLE_DEBEZIUM_JSON是 Source 侧的解析格式若在 KafkaSink中同时指定message_value_fields与该格式组合会直接抛出操作不支持异常。依赖部署格式模块中 Kafka 相关依赖均为 provided scope生产环境需确保seatunnel-format-compatible-connect-json的 jar 位于 SeaTunnel 的 plugin 目录格式插件目录运行时由 Kafka 连接器加载。事件时间语义行内event_time取自 Kafka 消息时间戳且仅在缺失时填充若上游消息时间戳非法 0则不设置对事件时间敏感的作业应结合 Kafka 消息的timestampType理解其含义。与 Debezium 格式的区分SeaTunnel 还有专门的DEBEZIUM_JSON格式见 KafkaSourceConfig支持debezium_record_include_schema与debezium_record_table_filter等选项。COMPATIBLE_KAFKA_CONNECT_JSON走的是 Kafka 官方JsonConverter通用路径更贴合 Connect JDBC 等标准 Connect Source 的输出两者均能处理 Debezium 产出的 Connect 格式数据可按上游 Converter 配置选择。小结COMPATIBLE_KAFKA_CONNECT_JSON是 SeaTunnel Kafka 连接器对接 Kafka Connect 生态的关键格式通过schema声明目标结构、通过key/value_converter_schema_enabled对齐上游 Converter 的 envelope 配置底层复用 Kafka 官方JsonConverter完成 Envelope 解析与 payload 提取并自动处理批量数组消息、事件时间挂载、表路由与消息头透传。配合 Jdbc 等 Sink 的generate_sink_sql能力即可快速搭建Kafka Connect 采集 → Kafka → 数仓/业务库的整条链路。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考