
Apache Kafka 中的 Kafka Connect统一数据集成框架的架构、特性与实践【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafkaKafka Connect 是 Apache Kafka 生态中面向数据集成的基础设施组件它以标准化的方式在 Kafka 与外部系统之间可靠、可扩展地搬运数据。本文以官方 Kafka Connect 概述文档 为核心结合 用户指南 与本仓库源码深入讲解 Connect 的定位、六大核心特性以及其底层实现与实战要点帮助你理解并上手这一连接一切数据系统的统一框架。Kafka Connect 是什么Kafka Connect 是一种在 Apache Kafka 与其他系统之间可扩展、可靠地流式传输数据的工具。它让开发者能够快速定义connector连接器将大批量数据搬入或搬出 Kafka数据摄入ingestion可以把整个数据库的数据、所有应用服务器上的指标采集进 Kafka topic供下游以低延迟进行流处理数据导出export可以把 Kafka topic 中的数据投递到二级存储、查询系统或批处理系统中用于离线分析。从仓库结构可以看到Kafka Connect 是一个完整的独立子项目connect/目录下包含 api连接器 SPI 接口、runtimeWorker 运行时与分布式协调、file示例文件连接器、mirror跨集群镜像以及 transforms消息转换等模块共同构成一套完整的集成框架。Kafka Connect 的六大核心特性官方概述明确了 Kafka Connect 的六个关键特性下面逐一展开并补充源码级佐证。1. 统一的连接器框架Kafka Connect 将其他数据系统与 Kafka 的集成方式标准化简化了连接器的开发、部署与管理。任何连接器只需实现统一的 SPI 接口即可被框架托管源连接器实现org.apache.kafka.connect.source.SourceConnector目标连接器实现org.apache.kafka.connect.sink.SinkConnector。以仓库自带的示例 FileStreamSourceConnector.java 为例它继承SourceConnector通过taskClass()返回任务类、taskConfigs(int maxTasks)产出任务配置、config()声明配置定义ConfigDef框架据此完成连接器的生命周期管理。开发者只需要关心从哪读、往哪写其余交给框架。2. 分布式与单机两种运行模式Connect 既能向上扩展为支撑整个组织的、集中管理的大型服务也能向下收缩为开发、测试和小规模生产部署。单机模式standalone所有工作在一个进程中完成配置简单、易于上手适合只需一个 worker 的场景例如采集日志文件但不具备容错能力。启动方式bin/connect-standalone.sh config/connect-standalone.properties [connector1.properties connector2.json …]分布式模式distributed自动负载均衡支持动态扩缩容且对运行中的任务、配置和 offset 提交数据都提供容错。启动方式bin/connect-distributed.sh config/connect-distributed.properties两者的差异体现在启动的类与配置上单机模式由 ConnectStandalone.java 驱动分布式模式由 ConnectDistributed.java 驱动。分布式模式下StandaloneHerder被 DistributedHerder.java 取代——后者在独立的线程中运行主循环一方面驱动组协调协议响应 rebalance 事件另一方面处理发往 leader 的外部请求如创建、修改、删除连接器。源码注释明确指出只有 leader 才能写配置 topic、分配任务follower 节点收到的写请求会被转发给 leader详见 DistributedHerder.java 中NotLeaderException的抛出逻辑。Worker 的通用必需配置单机与分布式都要设置配置项说明bootstrap.servers用于建立与 Kafka 集群初始连接的 broker 地址列表key.converter键的转换器类决定写入/读出 Kafka 的键序列化格式与连接器无关常见如 JSON、Avrovalue.converter值的转换器类决定值的序列化格式同样独立于连接器plugin.path默认null包含 Connect 插件连接器、转换器、转换的路径列表。注意示例中的FileStreamSourceConnector/FileStreamSinkConnector打包在connect-file产物中默认不包含在 CLASSPATH 或 plugin.path 中运行 quickstart 前必须把包含它的绝对路径加进来本仓库build.gradle中该产物名为connect-file具体 jar 文件名随版本变化当前仓库版本为4.5.0-SNAPSHOT单机模式专属配置配置项说明offset.storage.file.filename存储源连接器 offset 的文件仓库中的 connect-standalone.properties 给出了完整示例bootstrap.serverslocalhost:9092、key.converterorg.apache.kafka.connect.json.JsonConverter、value.converterorg.apache.kafka.connect.json.JsonConverter、key.converter.schemas.enabletrue、offset.storage.file.filename/tmp/connect.offsets以及便于调试的offset.flush.interval.ms10000。Kafka 客户端参数的继承与覆盖规则worker 配置中的参数面向 Connect 访问配置、offset、状态三个内部 topic 的 producer/consumer。源任务使用的 producer 与 sink 任务使用的 consumer 需要分别以producer.与consumer.前缀重复配置唯一无前缀继承的是bootstrap.servers大多数情况下同一个集群够用。例外是安全集群连接所需参数最多要在 worker 配置中设置三份管理访问、Kafka 源、Kafka sink。连接器级还可以用producer.override.与consumer.override.前缀做独立覆盖。3. REST 接口通过易用的 REST API 即可向 Connect 集群提交并管理连接器单机与分布式模式下均可用。REST 服务由listeners配置控制支持http/https默认监听8083 端口。核心端点包括GET /connectors列出活跃连接器POST /connectors创建连接器请求体为含name字符串与config对象的 JSON可选的initial_state支持STOPPED/PAUSED/RUNNINGGET|PUT|PATCH /connectors/{name}/config读取/更新/局部修改配置PATCH 中null值表示删除该键GET /connectors/{name}/status查看连接器及其任务状态、所在 worker、错误信息PUT /connectors/{name}/pause|stop|resume暂停、停止、恢复连接器stop 会释放任务占用的资源POST /connectors/{name}/restart?includeTaskstrue|falseonlyFailedtrue|false重启连接器与任务GET /connectors/{name}/topics与PUT /connectors/{name}/topics/reset查看/重置连接器活跃 topic 集合offset 管理端点GET /connectors/{name}/offsets、DELETE /connectors/{name}/offsets、PATCH /connectors/{name}/offsets后两者要求连接器处于 stopped 状态GET /connector-plugins、GET /connector-plugins/{plugin-type}/config、PUT /connector-plugins/{connector-type}/config/validate插件信息与配置校验GET /返回 Connect 集群基础信息worker 版本、git commit、Kafka 集群 IDGET|PUT /admin/loggers[/{name}]查看与动态调整日志级别需在admin.listeners上启用未配置时默认与普通 listeners 共用。注意分布式模式下连接器配置不通过命令行传递必须使用 REST API 创建、修改与销毁。集群中 follower 节点的某些 REST 请求会被转发到 leader若节点对外可达的 URI 与监听 URI 不同可用rest.advertised.host.name、rest.advertised.port、rest.advertised.listener调整。4. 自动 offset 管理只需连接器提供少量信息Kafka Connect 就能自动完成 offset 提交连接器开发者无需操心这一易错环节。offset 的存储位置取决于运行模式单机模式写入offset.storage.file.filename指定的文件分布式模式写入 Kafka 内部 topic。以文件源连接器为例FileStreamSourceTask 的 offset 以{filename: test.txt}为 partition、{position: 30}为 offset 结构记录的是文件字节流位置。分布式模式下需要预先规划好三个内部 topic建议手动创建以控制分区数与副本数不创建则自动创建并使用默认值配置项要求group.id集群唯一名称用于组建 Connect 集群组绝不能与 consumer group ID 冲突config.storage.topic存储连接器与任务配置应单分区、多副本、开启压缩compactionoffset.storage.topic存储 offset应多分区、多副本、开启压缩status.storage.topic存储任务状态可多分区、多副本、开启压缩仓库示例 connect-distributed.properties 中即为group.idconnect-cluster、offset.storage.topicconnect-offsets、config.storage.topicconnect-configs、status.storage.topicconnect-status示例为适配单 broker 将副本因子设为 1生产环境建议默认 3。5. 分布式与可扩展Kafka Connect 构建在 Kafka 已有的组管理协议之上向集群中添加更多 worker 即可水平扩展。每个 Connect worker 都是消费组中的一员group.id即组名通过 WorkerCoordinator 参与组协调DistributedHerder在 rebalance 后按分配结果启动/停止任务。这也解释了为什么 Connect 集群的group.id必须全局唯一——它直接复用了 Kafka 的消费组机制。6. 流式/批处理集成利用 Kafka 自身的能力Kafka Connect 是连接流式与批处理数据系统的理想桥梁它既可把实时流数据灌入 Kafka 供流处理引擎消费也可将 Kafka 中的流式数据导出到面向批处理的分析系统实现两类数据系统的无缝衔接。连接器的通用配置连接器配置本质是键值映射可通过 REST 请求的 JSON 载荷两种模式皆可或单机模式的 properties 文件传入。通用选项配置项说明name连接器唯一名称重复注册会失败connector.class连接器的 Java 类支持全限定名、类名或别名如FileStreamSinkConnector可写为FileStreamSinktasks.max该连接器最多创建的任务数连接器可因达不到该并行度而创建更少任务key.converter/value.converter可选覆盖 worker 层的默认转换器Sink 连接器必须且只能二选一设置输入来源topics逗号分隔的输入 topic 列表topics.regex匹配输入 topic 的 Java 正则表达式。轻量消息转换Transformations连接器可配置转换链对单条消息做轻量修改便于数据整形与事件路由。配置方式transforms转换别名列表顺序即应用顺序transforms.$alias.type转换的全限定类名transforms.$alias.$transformationSpecificConfig该转换的具体配置。官方文档给出了一个完整实战给文件源连接器加静态字段。先在connect-standalone.properties中将key.converter.schemas.enable与value.converter.schemas.enable从true改为false以使用无 schema 的 JSON然后配置namelocal-file-source connector.classFileStreamSource tasks.max1 filetest.txt topicconnect-test transformsMakeMap, InsertSource transforms.MakeMap.typeorg.apache.kafka.connect.transforms.HoistField$Value transforms.MakeMap.fieldline transforms.InsertSource.typeorg.apache.kafka.connect.transforms.InsertField$Value transforms.InsertSource.static.fielddata_source transforms.InsertSource.static.valuetest-file-source未加转换时消费connect-testtopic 的输出为foo bar hello world加上转换后输出变为{line:foo,data_source:test-file-source} {line:bar,data_source:test-file-source} {line:hello world,data_source:test-file-source}内置转换清单Cast字段/整体类型转换、DropHeaders、ExtractField、Filter配合谓词选择性过滤、Flatten、HeaderFrom、HoistField、InsertField、InsertHeader、MaskField、RegexRouter、ReplaceField、SetSchemaMetadata、TimestampConverter、TimestampRouter、ValueToKey。其实现位于 connect/transforms 模块。谓词Predicates转换可通过谓词做到仅对满足条件的消息生效与Filter组合可实现选择性过滤。所有转换都隐含predicate与negate两个配置negatetrue反转匹配结果。内置谓词TopicNameMatches匹配 topic 名符合指定正则的记录HasHeaderKey匹配含指定 header key 的记录RecordIsTombstone匹配 tombstone 记录值为 null。例如过滤掉footopic 的全部消息、并对除bar外所有 topic 应用ExtractFieldtransformsFilter,Extract transforms.Filter.typeorg.apache.kafka.connect.transforms.Filter transforms.Filter.predicateIsFoo transforms.Extract.typeorg.apache.kafka.connect.transforms.ExtractField$Key transforms.Extract.fieldother_field transforms.Extract.predicateIsBar transforms.Extract.negatetrue predicatesIsFoo,IsBar predicates.IsFoo.typeorg.apache.kafka.connect.transforms.predicates.TopicNameMatches predicates.IsFoo.patternfoo predicates.IsBar.typeorg.apache.kafka.connect.transforms.predicates.TopicNameMatches predicates.IsBar.patternbar错误处理与死信队列DLQ默认情况下转换/转换器阶段任何错误都会导致连接器失败fail-fast等价于如下默认配置# disable retries on failure errors.retry.timeout0 # do not log the error and their contexts errors.log.enablefalse # do not record errors in a dead letter queue topic errors.deadletterqueue.topic.name # Fail on first error errors.tolerancenone可调整这些配置实现重试、日志记录与死信队列投递例如# retry for at most 10 minutes times waiting up to 30 seconds between consecutive failures errors.retry.timeout600000 errors.retry.delay.max.ms30000 # log error context along with application logs, but do not include configs and messages errors.log.enabletrue errors.log.include.messagesfalse # produce error context into the Kafka topic errors.deadletterqueue.topic.namemy-connector-errors # Tolerate all errors. errors.toleranceall开启errors.log.include.messagestrue会连问题记录的 key/value/headers 一起写入日志可能泄露敏感信息需谨慎。精确一次Exactly-once语义Sink 连接器自 0.11.0 起需将 worker 属性consumer.isolation.level设为read_committed或通过连接器客户端配置覆盖策略设置consumer.override.isolation.levelread_committed使消费组忽略已中止事务中的记录源连接器自 3.3.0 起需在集群中启用框架级支持且仅在分布式模式下可用。新集群将exactly.once.source.support设为enabled已有集群需两轮滚动升级先设为preparing再设为enabled。需要强调的是精确一次能否达成高度依赖具体连接器的设计即便 worker 配置正确连接器若无法利用框架能力也无法实现。插件发现与性能权衡plugin.discovery配置控制 worker 发现插件类的策略对启动时间影响显著策略说明only_scan3.6 之前唯一行为兼容所有插件但较慢hybrid_warn3.6 起默认兼容所有插件但对不兼容service_load的插件打印 WARN 日志hybrid_fail发现不兼容插件即以错误停止 workerservice_load禁用旧式扫描改用更快的ServiceLoader机制不兼容的插件可能不可用升级到service_load前应验证兼容性用hybrid_warn启动并观察org.apache.kafka.connect包的 WARN 日志或先在测试环境用hybrid_fail验证启动。开发者为插件添加兼容性需在META-INF/services/下按超类类型如SinkConnector、SourceConnector、Converter、Transformation、Predicate等九类添加 ServiceLoader 清单文件每行一个全限定类名。安全考量Connect 允许运行任意代码的插件安装前必须信任插件来源REST API 默认不设防任何能访问它的人都能启停连接器只应向受信任用户开放否则极易在 worker 上获得任意代码执行能力自 Kafka 4.2.0 起建议将connector.client.config.override.policy设为AllowlistKafka 5.0.0 起成为默认值只放行确实需要覆盖的配置像sasl.jaas.config、sasl.login.class这类可加载类、可执行代码的配置仅在 REST API 仅对可信用户开放时才应放行。针对安全集群worker principal 与各连接器 principal 还需要按官方文档中的 ACL 表进行授权如 worker 需要对config.storage.topic、offset.storage.topic、status.storage.topic具备读写权限启用精确一次时还需connect-cluster-${groupId}事务 ID 的权限等完整 ACL 矩阵见 用户指南安全章节。快速上手路径快速体验 Kafka Connect先按 quickstart 启动本地 Kafka 集群与单机版 Connect再通过 REST API 或命令行参数提交连接器。官方共提供三份示例 worker 配置connect-standalone.properties、connect-distributed.properties、connect-mirror-maker.properties以及 connect-file 文件源/文件目标示例连接器源码见 connect/file可作为理解连接器 SPI 的最佳起点。完整的 REST API 规范可查阅本仓库 generator 生成的 OpenAPI 文档generator/目录与docs/apis/下的相关说明更深入的管理操作运行模式、REST 端点细节、错误报告、插件发现迁移可继续阅读 Kafka Connect 用户指南。【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考