十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

SeaTunnel Kafka 源连接器深度指南:配置解析、消息格式与生产实践

SeaTunnel Kafka 源连接器深度指南:配置解析、消息格式与生产实践 SeaTunnel Kafka 源连接器深度指南配置解析、消息格式与生产实践【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文以 Kafka.md 为骨架结合 connector-kafka 模块的源码实现系统讲解 SeaTunnel Kafka Source 的引擎支持、全部源选项、消息格式json / text / canal_json / debezium_json / avro / protobuf / NATIVE、消费位点控制、动态分区发现、认证安全SASL/SCRAM、IAM、Kerberos以及端到端精确一次等生产级用法。读者学完后可独立完成 Kafka 接入 SeaTunnel 的作业配置、故障排查与性能调优。功能总览Kafka 源连接器用于将 Apache Kafka 中的数据读取进 SeaTunnel 作业是 SeaTunnel 数据集成链路中使用最广泛的 Source 之一。它支持以下引擎Spark / Flink / Seatunnel Zeta从官方功能清单看连接器具备以下能力详见 connector-v2-features.md特性支持情况批处理✅ 支持流处理✅ 支持精确一次Exactly Once✅ 支持列投影❌ 不支持并行度✅ 支持用户定义拆分❌ 不支持从源码结构看KafkaSource.java 实现了SeaTunnelSource与SupportParallelism接口其getBoundedness()方法根据JobMode.BATCH返回BOUNDED有界、否则返回UNBOUNDED无界这正是批/流两种模式切换的底层依据。依赖安装使用 Kafka 连接器需要以下依赖可通过$SEATUNNEL_HOME/bin/install-plugin.sh下载或从 Maven 中央仓库获取org.apache.seatunnel:connector-kafkaKafka 通用版本。同时plugin-mapping.properties 中定义了插件注册关系seatunnel.source.Kafka connector-kafka seatunnel.sink.Kafka connector-kafka即 Source 与 Sink 共用同一个connector-kafka模块本文只讨论其 Source 侧。连通性 dry-runZetaSeaTunnel Zeta 的--dry-run connect子命令支持对 Kafka Source 做仅元数据的连通性校验不需要真正启动消费使用配置的bootstrap.servers和kafka.config中的安全设置通过describeTopics校验显式配置的主题当pattern true时使用与正常运行完全一致的主题全名匹配规则进行正则匹配同时支持tables_configs与旧版table_list两种多表配置输出 schema 复用正常 Source 的配置路径包括 native 字段、Kafka header 字段与事件时间元数据。其实现位于 KafkaSourceDryRunValidator.java关键设计如下30 秒时间预算常量MAX_TIMEOUT_MS 30_000元数据请求共享该预算若kafka.config.default.api.timeout.ms配置得更小则取较小值超时前置校验在应用 dry-run 限制前若显式配置的default.api.timeout.ms小于request.timeout.ms会抛出ConfigException——这与 Kafka 正常启动时的行为一致有界等待所有KafkaFuture均通过remainingMillis(deadline)计算剩余时间超时即抛TimeoutException客户端关闭也使用有界等待admin.close(Duration.ofMillis(1))避免线程阻塞不做任何副作用校验不会创建 consumer 或 producer不会读写消息、访问或提交消费位点也不会创建缺失的主题。需要特别注意的是校验成功仅证明可以访问元数据并不代表具备消费消息、访问消费组或反序列化实际消息的能力同时与运行时行为一致正则表达式当前没有可见的匹配主题时也视为合法只校验主题列表访问不验证未来主题的访问权限。目前 Kafkasink尚不支持连通性 dry-run。源选项详解以下为 Kafka.md 中定义的完整源选项表多数选项的定义可一一对应到 KafkaSourceOptions.java 与 KafkaBaseOptions.java 中的Option声明。名称类型是否必填默认值描述topicString否-使用表作为数据源时要读取数据的主题名称。未使用tables_configs或table_list时需要配置。支持逗号分隔多个主题例如topic-1,topic-2tables_configsList否-推荐使用的多主题表配置。当不同 topic 需要不同 schema 或 format 时使用。topic、tables_configs、table_list三者只能配置一个table_listList否-旧版兼容的多主题表配置。topic、tables_configs、table_list三者只能配置一个bootstrap.serversString是-逗号分隔的 Kafka brokers 列表patternBoolean否false若为true则使用指定的正则表达式匹配并订阅主题consumer.groupString否SeaTunnel-Consumer-GroupKafka 消费者组 ID用于区分不同消费者组commit_on_checkpointBoolean否true为 true 时仅在 SeaTunnel checkpoint 完成后提交消费者偏移量并禁用 Kafka 自动提交为 false 时禁用 checkpoint 提交并启用 Kafka 自动提交poll.timeoutLong否10000kafka 主动拉取时间间隔毫秒kafka.configMap否-除上述必要参数外可指定多个消费者客户端参数覆盖 Kafka 官方文档中定义的全部消费者参数schemaConfig否-数据结构包括字段名称和字段类型详见 Schema 特性formatString否json数据格式可选text、canal_json、debezium_json、ogg_json、maxwell_json、avro、protobuf和native。默认字段分隔符为,可通过field_delimiter自定义。canal 格式见 canal-jsondebezium 格式见 debezium-jsonavro_schemaString否-当format为avro时生效用于提供二进制 Avro 消息的 writer schema适用于 record 名称、namespace 或 union 结构与 SeaTunnel schema 不完全一致的场景format_error_handle_wayString否fail数据格式错误的处理方式可选fail与skip。fail时格式错误将阻塞并抛出异常skip时跳过该行数据debezium_record_include_schemaBoolean否true当format为debezium_json时生效说明 Debezium 记录中是否携带 schema 信息debezium_record_table_filterConfig否-用于过滤 debezium 格式的数据仅当格式为debezium_json时使用field_delimiterString否,自定义数据格式的字段分隔符start_modeStartMode否group_offsets消费者初始消费模式earliest、group_offsets、latest、specific_offsets、timestampstart_mode.offsetsConfig否-用于specific_offsets消费模式的偏移量start_mode.timestampLong否-用于timestamp消费模式的时间start_mode.end_timestampLong否-用于timestamp消费模式的结束时间仅支持批模式partition-discovery.interval-millisLong否-1动态发现主题和分区的间隔时间ignore_no_leader_partitionBoolean否false是否忽略没有 leader 的分区。为 true 时在分区发现过程中跳过无 leader 分区common-options-否-源插件常见参数详见 Source Common Optionsprotobuf_message_nameString否-格式为 protobuf 时有效指定消息名称protobuf_schemaString否-格式为 protobuf 时有效指定 Schema 定义strip_schema_registry_headerBoolean否false格式为 protobuf 或 avro 时有效。protobuf 在反序列化前去除 Confluent Schema Registry 头avro 去除固定的 5 字节头magic byte 与 schema ID。avro 启用时必须同时配置avro_schema且不会查询 Schema Registryreader_cache_queue_sizeInteger否2Fetcher 与 Reader 线程之间缓冲队列的容量每个元素是一次consumer.poll()的整批结果is_nativeBoolean否false支持保留 record 的源信息kafka_headers_fieldsArray否-指定要从 Kafka 消息 header 中提取并映射为行字段的 header key 列表每个 header 值以 STRING 类型追加到输出行末尾位于正常 schema 字段之后。不支持 NATIVE 格式重要从 checkpoint 或 savepoint 恢复时Kafka Source 会优先使用 checkpoint 中保存的 split offset。start_mode与 consumer group offset 只在首次启动或为尚未存在 checkpoint 状态的新发现分区初始化位点时生效。读取一个或多个 topic 时使用topic若不同 topic 需要不同的 schema 或 format请使用tables_configs。topic、tables_configs、table_list三者互斥不能同时配置。reader_cache_queue_size连接器在 Fetcher 线程和 Reader 线程之间缓冲 poll 结果批次的最大数量。在 KafkaSource.java 的createReader()中可以看到缓冲队列是一个容量为readerCacheQueueSize的LinkedBlockingQueueBlockingQueueRecordsWithSplitIdsConsumerRecordbyte[], byte[] elementsQueue new LinkedBlockingQueue(kafkaSourceConfig.getReaderCacheQueueSize());使用时需要注意队列中的每个元素是一次完整的consumer.poll()结果最多包含max.poll.records默认 500条消息当下游产生背压时队列可能被填满此时驻留在内存中的消息数上限约为reader_cache_queue_size × max.poll.records当消费的消息体较大时过高的值会导致堆内存占用过高。若观察到内存压力请减小此值或降低kafka.config中的max.poll.records。debezium_record_table_filter当format debezium_json时可用该配置按库、schema、表名过滤数据。配置示例如下debezium_record_table_filter { database_name test schema_name public // null 如果不存在 table_name products }此时只有test.public.products表的数据会被消费。从 KafkaSourceConfig.java 的实现看连接器会将三元组拼成TablePath再包装为DebeziumJsonDeserializationSchemaDispatcher进行分发过滤。元数据支持EventTimeKafka 源会在ConsumerRecord.timestamp 0时将其自动写入 SeaTunnel 行的EventTime元数据。这一行为的底层实现在 KafkaEventTimeDeserializationSchema.java它包装实际的DeserializationSchema在每次deserialize后将当前 record 的 timestamp 通过MetadataUtil.setEventTime(row, timestamp)附加到行上。借助 Metadata 转换 可以把这段时间戳暴露为普通字段方便做分区或下游 SQL 处理source { Kafka { plugin_output kafka_raw topic seatunnel_topic bootstrap.servers localhost:9092 format json } } transform { Metadata { plugin_input kafka_raw plugin_output kafka_with_meta metadata_fields { EventTime kafka_ts # ConsumerRecord.timestamp (ms) } } Sql { plugin_input kafka_with_meta plugin_output kafka_enriched query select *, FROM_UNIXTIME(kafka_ts/1000, yyyy-MM-dd, Asia/Shanghai) as pt from kafka_with_meta where kafka_ts 0 } }任务示例以下示例覆盖 Kafka Source 的主流使用场景。若尚未安装和部署 SeaTunnel请先参考 安装指南 进行安装部署再按照 快速开始 运行任务。简单示例多 topic 消费读取topic_1、topic_2、topic_3三个主题的数据并输出到客户端# 定义运行环境 env { parallelism 2 job.mode BATCH } source { Kafka { schema { fields { name string age int } } format text field_delimiter # topic topic_1,topic_2,topic_3 bootstrap.servers localhost:9092 kafka.config { client.id client_1 max.poll.records 500 auto.offset.reset earliest enable.auto.commit false } } } sink { Console {} }这里通过kafka.config直接透传 Kafka consumer 原生参数。注意enable.auto.commit false配合默认的commit_on_checkpoint true由 SeaTunnel 在 checkpoint 完成后统一提交偏移量。正则表达式主题pattern true时topic会被当作正则表达式匹配主题名source { Kafka { topic .*seatunnel*. pattern true bootstrap.servers localhost:9092 consumer.group seatunnel_group } }从 KafkaSourceSplitEnumerator.java 的实现看Enumerator 会通过adminClient.listTopics()获取集群全部主题再用pattern.matcher(topic).matches()做全名匹配非部分匹配因此正则必须完整匹配主题全名。动态发现分区流任务运行期间如果 Kafka topic 会新增分区可以配置partition-discovery.interval-millis定时发现新分区。新分区没有 checkpoint 位点时会按start_mode初始化消费位置env { job.mode STREAMING checkpoint.interval 5000 } source { Kafka { topic seatunnel_topic bootstrap.servers localhost:9092 consumer.group seatunnel_group start_mode latest partition-discovery.interval-millis 5000 format json } }源码层面当discoveryIntervalMillis 0时Enumerator 的open()会启动一个名为kafka-partition-dynamic-discovery的守护线程以scheduleWithFixedDelay周期执行discoverySplits()新发现的分区会进入 pending 队列并自动分配无需重启作业。注意该配置默认值为-1即默认关闭动态发现。读取 Kafka 消息 Header使用kafka_headers_fields将指定的 Kafka 消息 header 提取为行字段header 值以 STRING 类型追加在正常 schema 字段之后注意不支持NATIVE格式——该格式已经通过MapString, String字段暴露了所有 header。source { Kafka { topic my-topic bootstrap.servers localhost:9092 kafka_headers_fields [correlation-id, x-trace-id] schema { fields { user_id int name string } } format json } }输出行将包含user_idint、namestring、correlation-idstring、x-trace-idstring。如果某条消息中不存在对应的 header key则该字段值为null。实现层面KafkaSourceConfig.java 会通过extendCatalogTableWithHeaderFields把 header 字段以 STRING 类型追加到 CatalogTable 的列定义中并用KafkaHeadersDeserializationSchema包装原反序列化器完成 header 值注入。此功能与 Kafka sink 连接器的kafka_headers_fields对应支持 header 在 topic 间的全链路传递。NATIVE 格式保留 Kafka 原始信息如果需要保留 Kafka 原生的 key、partition、timestamp、headers 等信息可将format设为NATIVEsource { Kafka { topic test_topic_native_source bootstrap.servers kafkaCluster:9092 start_mode earliest format_error_handle_way skip format NATIVE value_converter_schema_enabled false consumer.group native_group } }返回的数据结构如下{ headers: { header1: header1, header2: header2 }, key: dGVzdF9ieXRlc19kYXRh, partition: 3, timestamp: 1672531200000, timestampType: CREATE_TIME, value: dGVzdF9ieXRlc19kYXRh }注意key/value是byte[]类型。从 KafkaSourceConfig.java 的nativeTableSchema()可以看到NATIVE 格式固定输出 7 列headersMapString,String、keybyte[]、offsetlong、partitionint、timestamplong、timestampTypestring、valuebyte[]。多 Kafka 源示例tables_configs当不同的 Kafka 主题需要不同的 schema 与 format 时应使用tables_configs官方提示Kafka 是非结构化数据源应使用tables_configs将来会删除table_list。以下示例根据主题与格式解析数据并基于 ID 执行 upsert 操作env { execution.parallelism 1 job.mode BATCH } source { Kafka { bootstrap.servers kafka_e2e:9092 tables_configs [ { topic ^test-ogg-sou.* pattern true consumer.group ogg_multi_group start_mode earliest schema { fields { id int name string description string weight string } }, format ogg_json }, { topic test-cdc_mds start_mode earliest schema { fields { id int name string description string weight string } }, format canal_json } ] } } sink { Jdbc { driver org.postgresql.Driver url jdbc:postgresql://postgresql:5432/test?loggerLevelOFF user test password test generate_sink_sql true database test table public.sink primary_keys [id] } }旧版兼容写法使用table_list结构与上述完全一致仅顶层键名不同env { execution.parallelism 1 job.mode BATCH } source { Kafka { bootstrap.servers kafka_e2e:9092 table_list [ { topic ^test-ogg-sou.* pattern true consumer.group ogg_multi_group start_mode earliest schema { fields { id int name string description string weight string } }, format ogg_json }, { topic test-cdc_mds start_mode earliest schema { fields { id int name string description string weight string } }, format canal_json } ] } } sink { Jdbc { driver org.postgresql.Driver url jdbc:postgresql://postgresql:5432/test?loggerLevelOFF user test password test generate_sink_sql true database test table public.sink primary_keys [id] } }从 KafkaSourceConfig.java 的createMapConsumerMetadata()可以看出连接器会按tables_configs或table_list逐条构造ConsumerMetadata每个表独立持有 topic、pattern、start_mode、schema 与 format再以TablePath为键组织为元数据 Map。Protobuf 配置format设置为protobuf时需同时配置protobuf_message_name与protobuf_schemasource { Kafka { topic test_protobuf_topic_fake_source format protobuf protobuf_message_name Person protobuf_schema syntax proto3; package org.apache.seatunnel.format.protobuf; option java_outer_classname ProtobufE2E; message Person { int32 c_int32 1; int64 c_int64 2; float c_float 3; double c_double 4; bool c_bool 5; string c_string 6; bytes c_bytes 7; message Address { string street 1; string city 2; string state 3; string zip 4; } Address address 8; mapstring, float attributes 9; repeated string phone_numbers 10; } bootstrap.servers kafkaCluster:9092 start_mode earliest plugin_output kafka_table } }Protobuf with Schema Registry wire format当消费使用 Confluent Schema Registry 编码的 Protobuf 消息时需要将strip_schema_registry_header设置为true。连接器会自动检测并删除 Schema Registry 格式头部magic byte、schema id 和 message indexes然后再反序列化 Protobuf 消息source { Kafka { topic test_protobuf_schema_registry_topic format protobuf strip_schema_registry_header true protobuf_message_name Person protobuf_schema syntax proto3; package org.apache.seatunnel.format.protobuf; option java_outer_classname ProtobufE2E; message Person { int32 c_int32 1; int64 c_int64 2; float c_float 3; double c_double 4; bool c_bool 5; string c_string 6; bytes c_bytes 7; message Address { string street 1; string city 2; string state 3; string zip 4; } Address address 8; mapstring, float attributes 9; repeated string phone_numbers 10; } bootstrap.servers kafkaCluster:9092 start_mode earliest plugin_output kafka_table } }注意启用strip_schema_registry_header时连接器可以安全地处理 Schema Registry 编码的消息和纯 Protobuf 消息。如果未检测到 Schema Registry 头部它会自动回退到标准 Protobuf 反序列化。从 KafkaSourceConfig.java 看该选项开启时会选用SchemaRegistryAwareProtobufDeserializationSchema。Avro 反序列化当 Avro 消息的 record 名称、namespace 或 union 结构与 SeaTunnel schema 不一致时需将format设置为avro并提供avro_schema若不提供avro_schema连接器会从用户配置的schema块派生解码 schema并同时作为 reader schema 与 writer schema 使用若生产端的 Avro 结构record 名称、namespace、union 结构与 SeaTunnel schema 不一致请显式配置avro_schema当前实现中没有Confluent Schema Registry 查询或按消息回退读取 schema 的机制。source { Kafka { topic users_avro bootstrap.servers localhost:9092 format avro avro_schema { type: record, name: User, namespace: com.example, fields: [ {name: id, type: long}, {name: name, type: string}, {name: email, type: [null, string], default: null} ] } schema { fields { id bigint name string email string } } } }对于由 ConfluentKafkaAvroSerializer写入的消息可设置strip_schema_registry_header true并提供avro_schema。连接器通过开头的 magic byte 检测并剥离固定的 5 字节线上格式头magic byte0加 4 字节 schema ID后再解码不会查询 Schema Registry该选项默认关闭关闭时原始 Avro 行为保持不变。需要留意开启该选项而未配置avro_schema时KafkaSourceConfig.java 会抛出SeaTunnelJsonFormatException错误码ILLEGAL_ARGUMENT。AWS MSK SASL/SCRAM将以下${username}和${password}替换为 AWS MSK 中的配置值source { Kafka { topic seatunnel bootstrap.servers xx.amazonaws.com.cn:9096,xxx.amazonaws.com.cn:9096,xxxx.amazonaws.com.cn:9096 consumer.group seatunnel_group kafka.config { security.protocolSASL_SSL sasl.mechanismSCRAM-SHA-512 sasl.jaas.configorg.apache.kafka.common.security.scram.ScramLoginModule required username\username\ password\password\; } } }AWS MSK IAM使用 IAM 认证时需要先从aws-msk-iam-auth项目下载aws-msk-iam-auth-1.1.5.jar并放到$SEATUNNEL_HOME/plugin/kafka/lib目录下。确保 IAM 策略中包含kafka-cluster:Connect等权限Effect: Allow, Action: [ kafka-cluster:Connect, kafka-cluster:AlterCluster, kafka-cluster:DescribeCluster ],源配置示例source { Kafka { topic seatunnel bootstrap.servers xx.amazonaws.com.cn:9098,xxx.amazonaws.com.cn:9098,xxxx.amazonaws.com.cn:9098 consumer.group seatunnel_group kafka.config { security.protocolSASL_SSL sasl.mechanismAWS_MSK_IAM sasl.jaas.configsoftware.amazon.msk.auth.iam.IAMLoginModule required; sasl.client.callback.handler.classsoftware.amazon.msk.auth.iam.IAMClientCallbackHandler } } }Kerberos 认证示例使用 Kerberos 前请在启动 SeaTunnel 之前设置 JVM 参数java.security.krb5.conf或更新/etc/krb5.conf中的默认krb5.conf。源配置示例source { Kafka { topic seatunnel bootstrap.servers 127.0.0.1:9092 consumer.group seatunnel_group kafka.config { security.protocolSASL_PLAINTEXT sasl.kerberos.service.namekafka sasl.mechanismGSSAPI sasl.jaas.configcom.sun.security.auth.module.Krb5LoginModule required \n useKeyTabtrue \n storeKeytrue \n keyTab\/path/to/xxx.keytab\ \n principal\userxxx.com\; } } }忽略无 Leader 分区当处理可能存在临时 leader 问题的 Kafka 集群时可以配置连接器忽略没有 leader 的分区source { Kafka { topic test_topic bootstrap.servers localhost:9092 consumer.group test_group ignore_no_leader_partition true start_mode earliest } }当ignore_no_leader_partition true时连接器将在分区发现过程中跳过任何没有 leader 的分区允许作业继续处理其他健康的分区。该逻辑位于 KafkaSourceSplitEnumerator.java 的getTopicInfo()中对partitionInfo.leader() null的分区做过滤并记录 warn 日志。配合动态分区发现与 EXACTLY_ONCE 下游的流式作业一种常见的长时间运行模式使用 Kafka 源读取数据并写入 Kafka sink配合 checkpoint 与semantics EXACTLY_ONCE实现端到端精确一次。开启partition-discovery.interval-millis后新增分区会被自动发现无需重启作业env { parallelism 2 job.mode STREAMING checkpoint.interval 10000 } source { Kafka { topic orders bootstrap.servers localhost:9092 consumer.group orders_consumer start_mode group_offsets commit_on_checkpoint true partition-discovery.interval-millis 30000 format json schema { fields { order_id bigint user_id bigint amount double } } } } sink { Kafka { topic orders_sink bootstrap.servers localhost:9092 format json semantics EXACTLY_ONCE transaction_prefix orders_sink_job partition_key_fields [order_id] } }同样的写法也适用于format debezium_json可以从 Kafka Connect sink 输出的 Debezium 变更事件中消费变更数据并转发到下游。常见问题start_mode各取值有什么区别start_mode行为earliest从每个分区最早可用的 offset 开始消费latest只消费任务启动后新产生的消息group_offsets从消费组已提交的 offset 恢复消费specific_offsets从每个分区指定的 offset 开始消费timestamp从指定时间戳处或其后的第一条消息开始消费任务中断后重启时使用group_offsets恢复需要从头重放全量数据时使用earliest。这五种模式定义在 StartMode.java 枚举中并由 KafkaSourceSplitEnumerator.java 的setPartitionStartOffset()分别通过OffsetSpec.earliest()、OffsetSpec.latest()、OffsetSpec.forTimestamp(...)、listConsumerGroupOffsets(...)与specificStartOffsets完成位点初始化。补充两点细节timestamp模式下可使用start_mode.end_timestamp指定结束时间仅支持批模式KafkaSourceConfig.java 在解析时会校验起止时间戳均不得晚于当前时间否则抛出IllegalArgumentExceptionspecific_offsets模式的偏移量配置格式为topic-partition - offset例如my-topic-0: 100源码中用最后一个-分隔主题名与分区号。如何按 Kafka 消息 key 过滤同一 topic 中的消息将format设为NATIVE可以把 Kafka 原始元数据包括key字段作为记录的一部分暴露出来再用 SQL Transform 保留所需 key 值的消息source { Kafka { topic events bootstrap.servers localhost:9092 format NATIVE consumer.group my-group } } transform { Sql { plugin_input kafka_source plugin_output filtered query SELECT * FROM kafka_source WHERE key expected_key_base64 } }注意NATIVE 格式中key字段为base64 编码的字节数组。Kafka Source 支持哪些消息格式支持json、text、canal_json、debezium_json、ogg_json、avro、protobuf和NATIVEMessageFormat.java 枚举中还包含maxwell_json、compatible_debezium_json与compatible_kafka_connect_json。当需要将 Kafka 元数据headers、key、partition、timestamp作为记录字段使用时选择NATIVE格式。format avro默认读取原始 Avro 二进制消息对于 ConfluentKafkaAvroSerializer写入的消息设置strip_schema_registry_header true并配置avro_schema即可详见上文 Avro 小节。如何配置 SASL/Kerberos 认证除使用kafka.config内嵌认证参数外也可以用kafka.*前缀属性直接写在 Source 块顶层source { Kafka { bootstrap.servers broker:9092 topic secure-topic consumer.group my-group kafka.security.protocol SASL_PLAINTEXT kafka.sasl.mechanism GSSAPI kafka.sasl.kerberos.service.name kafka kafka.sasl.jaas.config com.sun.security.auth.module.Krb5LoginModule required useKeyTabtrue keyTab/etc/kafka/kafka.keytab principaluserREALM.COM; } }这些kafka.*属性最终会汇入 consumer 配置即 KafkaBaseOptions.java 中的kafka.config处理路径因此与kafka.config块写法等价可按团队习惯选择。消费组 offset 是如何提交的SeaTunnel 在checkpoint 完成时向 Kafka 提交 offset。需在env块中通过checkpoint.interval开启 checkpoint。使用start_mode group_offsets重启任务时将从上次 checkpoint 提交的 offset 恢复消费。底层实现位于 KafkaSourceReader.java 的notifyCheckpointComplete()先判断commit_on_checkpoint默认true为false时直接跳过提交从checkpointOffsetMap取出该 checkpoint 对应的MapTopicPartition, OffsetAndMetadata通过KafkaSourceFetcherManager.commitOffsets(...)异步提交提交失败仅记录 warn 日志提交成功后清理该 checkpoint 及更早的所有待提交记录。与之配合checkpoint 快照阶段snapshotState会把每个 split 当前的startOffset记录为待提交位点从 checkpoint/savepoint 恢复时KafkaSource.java 的restoreEnumeratorEnumerator 会优先使用 checkpoint 中保存的 split offset而不是重新按start_mode定位。变更日志各版本针对 Kafka 连接器的功能变更、修复与行为调整请查阅 connector-kafka 变更日志。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表