Kafka Connect详解:大数据ETL的得力助手
这几年搞数据平台,我最大的感受就是:数据源越来越多,格式越来越杂,从MySQL、PostgreSQL到MongoDB,从REST API到各种消息队列,想把这些数据统一汇聚到一个地方做分析,靠人工写脚本一个个同步,不仅累死人,还容易出纰漏。直到我把Kafka Connect用起来之后,整个数据接入层的维护成本才真正降下来。这篇文章我想好好聊聊Kafka Connect,它不是那种花里胡哨的框架,而是实打实帮你解决“数据怎么稳定、高效地流进Kafka、再流出去”这件事的工具。无论你是刚接触大数据ETL的新手,还是已经在生产环境里跟各种同步任务搏斗过的老手,这篇内容应该都能给你一些参考。
Kafka Connect是Apache Kafka生态里的一个核心组件,专门用来做数据集成。它干的事情说白了就两件:把外部系统的数据搬进Kafka,这叫Source;把Kafka里的数据搬到外部系统,这叫Sink。听起来好像很简单,但真正用过之后你就会发现,它把多数据源接入、偏移量管理、分布式调度、容错这些脏活累活全包了,你只需要专注在“数据怎么转换”这一层就行。对于数据工程师来说,它就是ETL管道里那个最靠谱的数据搬运工。
文章后面我会先从整体设计思路讲起,然后拆解Source和Sink的底层机制,再结合我自己的实操经验把完整的落地步骤捋一遍,最后把我踩过的坑和排查心得一并整理出来。内容偏实战,尽量少讲虚的。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
1. 整体设计思路:为什么ETL要选Kafka Connect
1.1 从“脚本搬运”到“框架集成”的转变
我最早做数据同步的时候,用的还是写Python脚本、配crontab调度那一套。每接一个数据源,就要写一套“拉数据-清洗-推送到目标”的代码,每套代码要考虑重试、断点续传、脏数据处理这些逻辑。刚开始数据源少还能扛,等数据源一多,问题就全来了:有的脚本跑挂了没告警,有的数据源改了schema导致解析失败,有的同步延迟越来越大却不知道瓶颈在哪。整个数据接入层就是一座随时可能塌的“屎山”。
Kafka Connect之所以能成为ETL场景的得力助手,核心在于它把“数据搬运”这个通用问题抽象成了一层框架。你不需要关心任务怎么调度、偏移量怎么记录、并发度怎么分配,这些框架都帮你处理好了。你需要做的只是选一个合适的连接器(Connector),或者按接口规范写一个自定义连接器,然后声明“我要从哪里读数据、写到哪去、多久同步一次”,剩下的事情交给Connect集群。
1.2 Kafka在ETL链路里扮演的角色
在一个典型的大数据ETL架构里,Kafka本身是作为“数据总线”存在的。生产系统产生的数据先实时进入Kafka,然后由各种消费者去处理。Kafka Connect在这里起到了“桥梁”的作用:在数据入口处,Source Connector把外部数据源源不断送进Kafka Topic;在数据出口处,Sink Connector把Kafka里的数据投递到数仓、搜索引擎、缓存或者另一个消息队列。
这种架构最大的好处是解耦。上游业务系统不需要知道数据最终去了哪里,下游消费系统也不需要关心数据从哪里来。数据在Kafka这个中间层缓冲,天然支持削峰填谷,下游即使短暂抖动也不至于丢失数据。对于大数据平台来说,这种“统一接入、统一分发”的模式,比点对点的直连要稳健得多。
1.3 单机模式与分布式模式的取舍
Kafka Connect支持两种运行模式,单机模式和分布式模式。单机模式适合开发调试或者数据量很小的场景,配置文件里指定一个任务,启动一个进程就完事。但我一贯的建议是:哪怕数据量不大,也尽量用分布式模式跑。原因很简单,分布式模式能给你带来两个最实用的能力——水平扩展和故障转移。
在分布式模式下,你可以启动多个Connect Worker节点组成一个集群,连接器任务会由集群自动分配。某台机器挂了,任务会被转移到其他节点继续跑,不会因为单点故障就整个接入层瘫掉。所有偏移量、配置信息都存在Kafka内部Topic里(config.storage.topic、offset.storage.topic、status.storage.topic),天然具备持久化能力。这个设计理念和Kafka本身“日志即存储”的思路一脉相承,算是Kafka Connect一个很聪明的设计。
2. 核心机制拆解:Source与Sink的底层原理
2.1 Source Connector:数据怎么进Kafka
Source Connector负责从外部系统拉取数据,然后写入Kafka Topic。以常见的JDBC Source Connector为例,它通过轮询数据库表来捕获新增和变更的数据。底层配置里有两个关键参数:incrementing.column.name和timestamp.column.name。前者用自增主键来识别新增行,后者用时间戳字段来识别更新时间变化后的行。两个参数同时配置,就能做到“新增和更新都能捕获”。
这里有一个很多新手容易忽略的点:timestamp.column.name模式下,如果某条数据的更新时间字段值等于上一次轮询记录的最大时间值,有可能被重复读取。因为轮询逻辑是“大于上次记录的最大值”,而如果两条数据恰好落在同一毫秒,就可能出现漏一条或重一条的情况。所以生产环境里我倾向于用incrementing模式做主键增量同步,或者干脆用Debezium这类基于日志解析的CDC工具来做实时同步,后面我会细讲。
2.2 Sink Connector:数据怎么出Kafka
Sink Connector的逻辑相对简单——订阅一个或多个Kafka Topic,消费消息,然后把消息写入目标系统。以JDBC Sink Connector为例,它会根据Topic名称自动映射目标表,Topic的key作为主键,value作为数据行。配置insert.mode=upsert时可以实现“存在即更新、不存在即插入”的语义,很适合做宽表聚合写入。
但Sink方向有个典型的坑,就是“Topic分区数和目标系统并发度”的匹配关系。Kafka消息是按分区组织的,而Sink任务的并发度由tasks.max控制。如果你设置tasks.max=1,哪怕Topic有12个分区,也只有一个任务在消费所有分区的数据,吞吐量自然上不去。正确的做法是把tasks.max设为大于等于Topic分区数,让每个任务消费一个或多个分区的数据,充分发挥并行能力。
2.3 连接器的三大核心接口
Kafka Connect的连接器开发遵循一套标准接口,理解这套接口对排查问题非常有帮助。最上层是Connector接口,负责定义任务配置、校验配置参数、分发任务。真正干活的是Task接口,它才是跑数据同步逻辑的单元。还有一个Converter接口,负责把Kafka里的字节数据转换成对象,常用的是JsonConverter和AvroConverter。
可以这么理解:Connector像是一个“包工头”,负责接活、分配活;Task像“工人”,每人都有一份清单,按清单执行具体的搬砖动作。框架层面通过REST API(/connectors、/connectors/{name}/tasks等端点)来管理这些连接器的生命周期。如果你要自定义连接器,动手前最好把这些接口的先决逻辑理清楚,否则写出来的东西很容易在task状态管理上出问题。
3. 实操落地:从零搭建一套Kafka Connect同步管道
3.1 环境准备与安装部署
我拿一套典型的测试环境举例:三台服务器组成的Kafka集群(版本2.8.0),Kafka Connect也以分布式模式跑在这三台机器上。不用额外装什么重组件,因为Kafka发行包本身自带connect-distributed.sh脚本,这就是分布式模式启动入口。
如果你是自己编译或者用Confluent Platform发行版,那更简单,Confluent的Hub上有大量现成连接器可以下载。这里我建议直接使用Confluent发行版里的连接器生态,因为开源Apache Kafka自带的连接器类型比较有限,主要是FileStreamSource和FileStreamSink这种演示性质的。真实场景你需要去Confluent Hub搜jdbc、elasticsearch、hdfs、mongo这些关键词,下载对应连接器jar包放到share/java/kafka-connect/目录下。
启动之前有三件事必须确认:bootstrap.servers配置指向Kafka集群地址;group.id每个集群要唯一;三个存储Topic(配置、偏移量、状态)如果不存在,Connect会自动创建但你需要确认自动创建Topic的开关是打开的。我第一次部署就吃过亏,Kafka集群配置了auto.create.topics.enable=false,结果Connect起不来,日志里全是一堆“topic not found”的错误。
3.2 用JDBC Source把MySQL数据导入Kafka
假设我们要把MySQL里的orders表实时同步到Kafka,最简单的做法是配置一个JdbcSourceConnector。下面这份配置我实测过,可以直接抄:
json复制{
"name": "mysql-orders-source",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
"connection.url": "jdbc:mysql://10.0.0.5:3306/business?useSSL=false",
"connection.user": "kafka_connect",
"connection.password": "xxxxxx",
"table.whitelist": "orders",
"mode": "incrementing",
"incrementing.column.name": "id",
"topic.prefix": "mysql-biz-",
"tasks.max": "1",
"poll.interval.ms": "5000"
}
}
这份配置跑起来后,效果就是每5秒轮询一次orders表,把新插入的行以Topic名mysql-biz-orders发送到Kafka。注意topic.prefix的作用,它和表名拼在一起才是最终Topic名,不要漏掉最后那个连字符,否则Topic名会变得很难看。
有人可能会问,为什么tasks.max只设为1?因为单表单库的轮询同步,JDBC Source的并发提升空间不大,一个任务就够了。但如果你要同步几十张表,那tasks.max可以调大一些,让多个任务并行去轮询不同表,吞吐量会明显提升。
3.3 用JDBC Sink把Kafka数据写回PostgreSQL
接下来把Kafka里的mysql-biz-orders数据同步到PostgreSQL的一张宽表里。配置如下:
json复制{
"name": "pg-orders-sink",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
"connection.url": "jdbc:postgresql://10.0.0.9:5432/analytics",
"connection.user": "sink_user",
"connection.password": "xxxxxx",
"topics": "mysql-biz-orders",
"insert.mode": "upsert",
"pk.fields": "id",
"auto.create": "true",
"auto.evolve": "true",
"tasks.max": "3"
}
}
这里auto.create设为true意思是目标表如果不存在,连接器会根据消息的schema自动建表。auto.evolve设为true意味着源表加了列,目标表会自动加列。开发环境非常方便,生产环境我不建议直接开这两个开关,最好由DBA审核后手动执行DDL,否则一个线上环境突然自动加了列,后续的权限审计和表结构管理会乱套。
运行起来之后,你会看到Topic里的每条消息都被消费并写入PostgreSQL对应表。如果数据量较大,记得把batch.size和max.retries这些参数根据业务容忍度做一些调整。
3.4 使用REST API管理连接器生命周期
Kafka Connect最方便的一点是有完整的REST API,所有操作都可以用命令行工具或者脚本控制。我常用的几个接口分享给大家:
bash复制# 查看集群里所有连接器
curl -s http://localhost:8083/connectors | jq
# 查看某个连接器的详细配置和状态
curl -s http://localhost:8083/connectors/mysql-orders-source/status | jq
# 暂停一个连接器(保留配置,停止拉数据)
curl -s -X PUT http://localhost:8083/connectors/mysql-orders-source/pause
# 恢复一个被暂停的连接器
curl -s -X PUT http://localhost:8083/connectors/mysql-orders-source/resume
# 删除连接器(会停止任务并删除配置)
curl -s -X DELETE http://localhost:8083/connectors/mysql-orders-source
这个API可以说是Kafka Connect最实用的“遥控器”。有一次线上数据源要做维护,我直接调pause接口把同步停了,等维护结束后再调resume恢复,整个过程不用重启任何Worker进程,数据管道一点没受影响。我强烈建议把这几个命令存成一个小脚本放在运维工具库里,后续排查问题会非常省事。
4. 转换器与数据格式:ETL链条上的关键决策点
4.1 JsonConverter vs AvroConverter vs SchemaConverter
数据进了Kafka,以什么格式存储、传输,是很多人一开始不会太在意、但后期很难改的决策。Kafka Connect里常见的转换器有三种:JsonConverter把消息以JSON字符串格式写入Kafka,直观、调试方便,但占用的存储空间大,且不强制schema一致性;AvroConverter配合Schema Registry使用,压缩率高、schema演进能力强,适合大规模生产环境;还有基于Protobuf的转换器,适合有Protobuf技术积累的团队。
我个人对这个问题的建议是:如果你只是个人学习、做Demo,JSON完全够用;如果数据量大、模型经常演进、要保证多团队协作时的字段规范,那就老老实实上Avro配Schema Registry。
这里需要插一个重要背景:Apache Kafka 3.0之后默认移除了对旧消息格式MessageFormat v0和v1的兼容支持,同时kafka-console-producer的默认序列化器也改了,所以如果你还在用老版本的习惯去生产数据给Connect消费,很可能会遇到“反序列化失败”的问题。用Avro或者JSON的时候,一定记得在Connect的配置文件里把key.converter和value.converter同时设置,不要一个配了另一个没配,那会死得很惨。
4.2 在Connect里做轻量级数据转换
我们很多时候只想改字段名、去掉一个嵌套字段,没必要专门写一套流处理程序。Kafka Connect提供了Single Message Transform(SMT)机制,可以在消息进入Kafka之前或写出Kafka之后做一些字段级操作。
我最常用的三个SMT是RenameField、TimestampConverter和InsertField。比如从MySQL同步过来的字段名是create_time,但目标数仓规范要求字段名是created_at,直接在Sink连接器里配一条转换规则就行:
json复制"transforms": "rename",
"transforms.rename.type": "org.apache.kafka.connect.transforms.RenameField$Value",
"transforms.rename.renames": "create_time:created_at"
SMT的定位是“轻量级”,复杂的数据清洗(比如多流join、聚合计算)还是应该交给Kafka Streams或Flink去处理,不要在一个Sink配置里堆几十个转换步骤,那样维护难度会急剧上升。记住一个原则——连接器里只做简单整形,一切有状态的计算都放到下游处理引擎里。
5. 高并发与数据一致性:生产环境绕不开的两个话题
5.1 水平扩展能力分析
你可能会担心,Kafka Connect能不能扛住突发流量。答案是可以,但并发度不是白白来的,它由tasks.max和底层系统能力共同决定。
比如同样一个JDBC Source连接器,从tasks.max=1提高到tasks.max=4,会发现轮询频率和拉取吞吐都会成倍提升。原理就是框架把这个连接器的一个任务拆分成多个Task,每个Task负责不同的表或者不同的分区集合。对Sink连接器来说,tasks.max通常应该与Topic分区数匹配,如果你有12个分区,那tasks.max配置成12能最大化消费并行度。配置成6也行,每个任务会消费两个分区,效果是吞吐量减半但资源占用更少。具体设多少,我一般看目标系统的写入能力,不让下游被打爆,也不让上游Kafka堆积太多。
连接器内部还有一层缓冲机制,比如JDBC Sink的batch.size和linger.ms就决定了攒多少条消息才批量写一次。调大batch.size可以降低下游数据库的写入次数,减少连接开销,但也会带来一定延迟。我的习惯是默认batch.size=3000,如果下游数据库出现锁等待或者死锁,再适当调小。
5.2 死信队列与错误处理策略
真实环境里,数据不会一直干干净净。源库突然写入了一个超出目标字段长度的字符串、枚举值对不上、或者序列化失败,这些异常都会导致Sink任务抛错。默认情况下,如果某条消息一直处理失败,Kafka Connect会进入一种“反复重试,直到成功或任务暂停”的状态,相当于一条脏数据卡住整条管道。
解决办法是启用死信队列(Dead Letter Queue,简称DLQ)。在配置里加上下面这条:
json复制"errors.tolerance": "all",
"errors.deadletterqueue.topic.name": "sink-dlq-orders",
"errors.deadletterqueue.context.headers.enable": "true"
这样处理失败的消息会被写入kafka-connect-dlq这个Topic,而不是把整个任务折磨到罢工。配置errors.deadletterqueue.context.headers.enable=true还能把失败原因、原始消息头部信息封装成header,后续排查时直接看DLQ Topic里的消息就知道是什么原因导致的失败。
这个机制救过我很多次。曾经有次业务方在源库某个字段里写入了包含非法UTF-8序列的内容,导致下游ES索引失败。如果没有DLQ,整条数据管道就停在那里了;有了DLQ,问题消息被隔离到专门的Topic里,主链路照常跑,等业务方修好数据后我们再手动把DLQ里的消息重新投递回去。
5.3 至少一次语义 vs 精确一次语义
Kafka Connect默认提供的是“至少一次”(at-least-once)投递语义。简单来说就是:数据可能会重复,但绝不会丢。因为Connect在写入目标系统成功之后才提交偏移量,如果写入成功后、偏移量提交前进程崩溃了,重启后会从旧偏移量重新消费,就会产生重复数据。这个语义在实际生产里是大多数场景都能接受的,下游做一层去重就行。
如果要追求“精确一次”,官方方案是Exactly Once特性,主要通过事务和幂等写入来实现。实际使用时要开启delivery.guarantee.exactly.once配置,并且确保目标系统是支持事务的,比如Kafka本身或者某些数据库。这里我不建议一开始就追求精确一次,它需要额外的配置与系统支持,还会引入性能开销。先做好去重逻辑,比在框架层纠结精确一次要划算得多。
6. 部署与运维:那些年我踩过的坑
6.1 配置文件的秘密:不要全堆在worker.properties里
很多初学者会把所有连接器的配置都写在worker.properties里,这样看起来“集中管理”,实际上非常难受。一旦连接器数量超过十个,这个文件就会变成一团乱麻,而且每次改一个连接器的配置都要重启整个Connect进程,影响面太大。
正确的做法是,worker.properties里只放基础配置,比如Kafka连接信息、转换器、偏移量存储Topic这些。具体的连接器配置通过REST API单独提交。这样你增删连接器、修改配置都不需要重启Worker,Kafka Connect会在后台动态加载新配置并重新调度任务。这一点在生产环境极其重要——如果某个连接器需要临时调参,你可以做到零停机更新。
6.2 日志排错:先看Connector日志,再看Worker日志
Kafka Connect的日志分成两层。第一层是Connect Worker进程的日志,记录的是框架级别的信息;第二层是每个连接器运行时产生的日志,往往打印在同一个worker日志中,但会用连接器名称和任务ID做标识。
遇到连接器启动失败,我的排查顺序是:
- 通过REST API查看任务状态,看是
RUNNING还是FAILED; - 用
curl获取任务实例日志,命令是curl -s http://localhost:8083/connectors/{name}/tasks,然后根据task id去日志文件里搜对应的异常堆栈; - 重点看有没有
Caused by,这通常是底层真实原因。
举一个我印象深刻的例子:一个Elasticsearch Sink Connector启动就报错,提示Could not find class ...。乍看以为是jar包没放好,后来一查日志才发现在worker升级之后,share/java/kafka-connect/目录下多了一个旧版本jar包,和新的jar包产生了类冲突。排查classpath问题的时候,不要只想“缺什么”,还要想“多了什么”。
6.3 关于偏移量丢失的记录
Kafka Connect的偏移量存放在内部Topic里。如果你不小心删除了offset.storage.topic,那么连接器会丢失偏移量记录,从“头”开始消费或者从源系统初始位置重新拉取数据。这个事故一旦发生,轻则数据重复一大片,重则直接把下游系统写爆。
所以生产环境一定要给内部Topic设置合理的副本数和保留策略。offset.storage.topic这种配置需要设置cleanup.policy=compact,确保偏移量不会因为过期被删除。不要因为这些Topic是自动创建的就不去管它们,默认的segment.bytes和retention.bytes可能完全不适合生产场景。
6.4 监控与告警
Kafka Connect暴露了JMX指标,可以接入Prometheus和Grafana。最值得监控的指标是:
kafka.connect:type=connect-metrics下的任务状态,状态不是RUNNING的都要告警;source-record-write-rate和sink-record-read-rate,观察吞吐量是否正常;- 各连接器的
offset-commit-failure-rate,如果升高说明偏移量提交有问题,需要赶紧查。
另外我强烈建议配置一个定期扫描REST API状态的脚本,比如每5分钟检查一次/connectors的状态,发现FAILED就发钉钉通知。因为连接器任务失败并不一定会导致整个进程退出,有时候就是默默地卡在那里,如果没有主动监控,你可能会在业务方投诉之后才发现数据已经断了好几个小时。
7. 常见问题与排查技巧实录
7.1 任务一直显示FAILED,但是日志没有完整堆栈
遇到这个问题先别慌,先确认是否开启了log4j的DEBUG级别。Kafka Connect的有些异常信息只在DEBUG级别下才会完整打印。另外很常见的原因是转型器(Converter)配置错误,尤其是key.converter和value.converter不一致,导致任务启动时就抛SerializationException。检查完配置之后,可以用kafka-console-consumer订阅相关Topic,看看里面消息数据是否正常,然后手动反序列化一条消息测试。
7.2 同步延迟越来越大,怎么办
同步延迟增大,一般有两个方向排查:
第一是源端拉取能力不足。如果是JDBC Source,先看poll.interval.ms是不是设得太大,轮询间隔5秒变成了5分钟,那延迟肯定高。把轮询间隔调小,再把batch.max.rows调大,可以让每次轮询拉取更多的数据。
第二是Sink端写入瓶颈。检查下游数据库的活跃连接数、锁等待时间,还有Sink任务的数量。如果目标表有大量索引,写入肯定快不起来。合理的表象是:Connect集群吞吐量稳定,下游数据库的QPS接近正常水位,Topic的消费Lag基本为0。
7.3 连接器配置更新了但没生效
这种问题我见得多了。记住,Kafka Connect的配置更新是“异步生效”的。你通过REST API提交了新的配置,返回200,但连接器可能需要几秒甚至几十秒才会重新加载配置并重启任务。如果你马上通过/status接口查看,可能还显示旧状态。千万不要手动重启Worker进程去“加速”,该做的只是等待,或者通过/connectors/{name}/tasks接口强行重启某个任务。
另外要特别注意:修改连接器配置时不要同时修改tasks.max,否则可能触发重新分配任务,让一批任务全部重启。你要是一次改完之后状态全乱了,最稳妥的办法是暂停连接器、更新配置、再恢复连接器。
7.4 Source端删除的数据,Sink端不会同步删除
这是JDBC轮询模式的一个天然限制。incrementing和timestamp模式只能捕获新增和更新,无法捕获删除操作。如果你的业务场景必须要求数据删除也能同步,那就得换基于CDC日志的方式,比如用Debezium的PostgresConnector或者MySQL的binlog连接器。这类连接器会把删除事件也封装成一条消息,下游通过识别消息里的op字段来执行删除。
所以在选型阶段,务必想清楚:你的数据同步需要的是“增量备份”级别的同步,还是“实时复制”级别的同步。前者用JDBC轮询就够了,后者上CDC工具更合适。选错了,后面返工的成本非常高。
8. 从J2SE到大数据协作:聊聊Kafka Connect周边生态的整合
写到这里有人可能会疑惑,为什么聊Kafka Connect要扯到J2SE这些基础能力。其实是因为我见过太多数据工程师在处理Connect链路时,遇到报错连日志都不看,最后发现问题的根源其实特别简单,比如JVM内存不足、GC停顿影响吞吐、classpath冲突等等。Kafka Connect本质上是跑在JVM上的分布式应用,你对JVM的调优经验完全可以迁移过来。
比如Connect Worker默认的堆内存可能只有1GB,当你挂了10个连接器、每个连接器又启动多个Task之后,内存很容易不够用,就会频繁GC甚至OOM。我的经验是:一个Worker节点上连接器数量较多时,堆内存至少给到4GB至8GB,同时调整-Xms和-Xmx为相同值,避免动态扩容时引入额外性能开销。
类加载问题也一样。安装第三方连接器时,统一把jar包放到share/java/kafka-connect/目录下,不要去改Worker类路径里已有的Kafka核心jar包。如果你既用了Confluent的连接器,又自己写了一个自定义连接器,注意两者的依赖版本冲突。实际操作中我遇到过Kafka客户端版本不一致导致RPC通信异常的案例,最终是统一所有连接器的Kafka客户端库版本才解决。
还有一点值得提醒:连接器插件最好使用隔离的classloader,Kafka Connect 2.3版本之后默认开启了插件路径隔离,但如果你用的发行版版本较老,需要确认plugin.path配置正确。把这些基础层面的问题搞定了,后面真正做数据管道开发的时候会顺利很多。
9. 实操总结与个人经验心得
如果只让我用一句话概括Kafka Connect,那就是:用标准化框架替代手工脚本,用配置管理替代代码管理,用可观测指标替代“跑完才知道成功与否”的黑盒操作。它的核心优势不在于性能压榨到极致,而在于稳定性和可维护性。数据量到了TB级别之后,真正卡你的往往不是同步速度,而是“管道能不能7x24小时不中断地跑完”。
个人实操中有几条经验想特别分享出来:
第一,连接器配置要纳入版本管理,不要只在服务器上改。用Git管理配置JSON,每次变更都走评审,出了问题能快速diff出哪一行配置变了,这能省下大量排查时间。
第二,连接器命名要规范。我见过有人把连接器命名为test1、test123,半年之后根本不知道哪个连接器对应哪个业务链路。我的规范是“目标系统-源系统-数据主题”,比如mysql-orders-to-pg、mongo-users-to-es,一眼就知道这个连接器在干什么。
第三,一定要做全链路的数据校验。Connect跑通了不代表数据是对的。我会定期对比源库某些表的行数、SUM值、或者抽样几条数据,和目标数仓里的数据做比对。Kafka Connect提供了tasks级别的状态信息,但数据正确性只能靠业务层校验闭环。
最后再分享一个个人心得。不要一上来就追求最复杂的架构,不要一上来就上CDC、上Schema Registry、上精确一次。先把一套简单的JDBC同步管道跑稳,理解连接器状态的流转过程,再把SMT、DLQ、监控加进去。数据接入层最怕的不是慢,而是“不可控”。Kafka Connect恰恰是那个把“不可控”变成“可控”的框架。用好了它,你的数据管道就能稳稳当当地支撑起整个大数据平台的上层应用。
