
Kafka connect 就是一个数据传输的工具可以使用connect非常方便的将外部系统当中的数据可以和Kafka之间做数据的传递。比如说文件比如说数据库。可以将文件中的数据导入到Kafka当中也可以将MySQL中的数据导入到Kafka中。还可以将Kafka中的数据导出到MySQLHDFS离线的本地文件等等这些外部的存储设备。比如也可以使用flume去采集日志文件的数据导入到Kafka中或者将Kafka的数据导出到HDFS这些其实也是可以去做到的。但是Kafka connect有它自己的优势比如Kafka connect可伸缩安全可靠而且是一个流式的处理。自动offset管理生产者生产的消息是有offset消费消费的消息也是有offset。如果要将数据库中的数据或者文件中的数据导入到Kafka那数据库和文件扮演的是生产者的角色。Kafka拉取数据得知道拉取到什么位置了。这个文件中有新的行出现的时候它得知道拉取什么样的数据。而kafka offset是自动维护的不需要手动去进行管理。Kafka connect 就是一个与外部存储系统数据系统之间的数据导入和导出的一个工具。角色角色本质核心职责Worker操作系统进程 (JVM)负责提供运行环境如 REST API、与 Kafka 通信、管理配置可以单节点运行Standalone也可以多节点组成集群Distributed。Connector逻辑配置与协调者负责定义**“从哪里读/写到哪里”配置管理。它本身不参与实际的数据传输**而是根据配置将整体工作切分为多个独立作业并创建 Task。Task实际干活的线程 (Worker 内)真正负责**数据抽取Source或数据写入Sink**的执行单元。Task 之间相互独立、无状态可以并行运行。Connector 划分工作当你向 Connect 提交一个配置比如读取一个包含 10 个 Partition 的 Kafka Topic 写入数据库Connector 决定将这个任务拆分为 3 个 Task假设tasks.max3。Worker 调度 Task单机模式 (Standalone)所有的 Connector 和 Task 都运行在同一个 Worker 进程内。分布式模式 (Distributed)集群中的多个 Worker 节点会自动打散并均衡分配这些 Task。如果某个 Worker 节点宕机集群会自动把该节点上的 Task 转移到其他 Worker 上继续运行Failover。Task 独立并行执行每个 Task 负责一部分数据的读取/写入例如 Task 1 读 Partition 0-3Task 2 读 Partition 4-6。它们在 Worker 内作为独立的线程高效运行。Kafka connect 单机模式启动命令至少两个参数一个是connect-stand-alone的一个配置从第二个参数往后就是各种connect的配置了。source是一种connectsink也是一种connect。所以要将source的配置和sink的配置都放在后面的配置文件当中。目录结构说明kafka conf目录下的文件 ├── connect-standalone.properties # Connect 运行进程的主配置文件 (配置 JSON Converter) ├── connect-file-source.properties # FileSource Connector 配置文件读文件 - Kafka ├── connect-file-sink.properties # FileSink Connector 配置文件Kafka - 写文件文件一connect-standalone.propertiesConnect 运行进程的主配置文件bootstrap.servers10.202.13.111:9092,10.202.12.107:9092,10.202.13.112:9092这两个是转换器当Kafka从外部的存储设备中读到了数据之后把这每一条的数据转换成什么格式在这里转成了json格式。Key 和 Value 的序列化方式使用 JSON 格式key.converterorg.apache.kafka.connect.json.JsonConvertervalue.converterorg.apache.kafka.connect.json.JsonConverter默认是打开的状态打开之后它能够有很多的键值对都存储到Kafka的topic当中。key.converter.schemas.enabletruevalue.converter.schemas.enabletrueKafka connect会自动的保存offset在stand alone的模式下这个offset是保存在一个文件当中的。默认是保存在这个文件当中。Standalone 模式下保存 Connector 偏移量Offset的文件路径offset.storage.file.filename/tmp/connect.offsets这个是自动保存offset一个偏移量一个时间10s。offset.flush.interval.ms10000上面文件的内容其实不需要去修改什么。connect-standalone.properties 内容bootstrap.servers10.20.13.111:9092,10.20.12.107:9092,10.20.13.112:9092 key.converterorg.apache.kafka.connect.json.JsonConverter value.converterorg.apache.kafka.connect.json.JsonConverter key.converter.schemas.enabletrue value.converter.schemas.enabletrue offset.storage.file.filename/tmp/connect.offsets offset.flush.interval.ms10000 plugin.path/apps/svr/kafka/libs errors.toleranceall errors.log.enabletrue errors.log.include.messagestrue文件二connect-file-source.properties读取本地 JSON 文件并发送到 Kafka Topic 的 Connector 配置。namelocal-file-source 这个可以随便取这个就是source的名字connector.classorg.apache.kafka.connect.file.FileStreamSourceConnector 这个是读取文件的类tasks.max1 这个任务要跑在多少个task当中stand alone只有一个taskfile/home/apps/access.log 要读取文件的位置topicconnect-test 读到了文件之后将消息存在哪个主题当中namelocal-file-source connector.classorg.apache.kafka.connect.file.FileStreamSourceConnector tasks.max1 file/home/apps/access.log topicconnect-test[appsTLKYVM202012107 kafka]$ bin/kafka-topics.sh --create --topic connect-test --bootstrap-server 10.202.13.111:9092,10.202.12.107:9092,10.202.1 3.112:9092 --partitions 3 --replication-factor 3 Created topic connect-test.文件三connect-file-sink.propertiesnamelocal-file-sinkconnector.classorg.apache.kafka.connect.file.FileStreamSinkConnectortasks.max1filetest.sink.txt 最后数据要保存到什么文件当中topicsconnect-testnamelocal-file-sink connector.classorg.apache.kafka.connect.file.FileStreamSinkConnector tasks.max1 file/home/apps/access.log.sink topicsconnect-test[apps kafka]$ ./bin/connect-standalone.sh -daemon config/connect-standalone.properties config/connect-file-source.properties config/connect-file-sink.properties[apps kafka]$ jps -l 58531 kafka.Kafka 9446 org.apache.zookeeper.server.quorum.QuorumPeerMain 84262 sun.tools.jps.Jps 20844 org.apache.kafka.connect.cli.ConnectStandalone数据未写入文件FileSink 无输出诊断通常是因为Topic 中尚无新消息、Connector 无法读取文件权限/路径错误、或者Converter 模式不匹配导致解析失败。1.检查 Source 文件的追加写入验证消息触发。FileStreamSourceConnector默认只读取启动后新增追加Append的内容或者从上次 offset 继续读。如果test-input.txt在启动前就已经写入完毕Connector 不会自动重新读取整份文件。在控制台运行以下命令追加一条新的 JSON 数据[apps ~]$ echo {id: 1, name: Kafka Test} access.log验证方式观察控制台日志是否有数据发送记录并检查access.log.sink是否产生数据。2.检查 Kafka Topic 是否有数据确认 Source 是否成功发消息。使用 Kafka 自带的命令行消费者确认FileStreamSource是否成功将数据写入到了 Kafka Topic 中[apps kafka]$ bin/kafka-console-consumer.sh --bootstrap-server 10.202.13.111:9092,10.202.12.107:9092,10.202.13.112:9092 --topic connect-test --from-beginning {schema:{type:string,optional:false},payload:{\id\: 1, \name\: \Kafka Test\}} {schema:{type:string,optional:false},payload:{\id\: 1, \name\: \Alice\, \role\: \admin\}} {schema:{type:string,optional:false},payload:{\id\: 1, \name\: \Alice\, \role\: \admin\}} {schema:{type:string,optional:false},payload:{\id\: 1, \name\: \Alice\, \role\: \admin\}} {schema:{type:string,optional:false},payload:{\id\: 1, \name\: \Alice\, \role\: \admin\}} {schema:{type:string,optional:false},payload:{\id\: 1, \name\: \Alice\, \role\: \admin\}} {schema:{type:string,optional:false},payload:{\message\: \hello kafka\}}3.排查 JsonConverter 格式匹配问题检查 Schema 开关。如果connect-standalone.properties中设置了value.converterorg.apache.kafka.connect.json.JsonConverter请核对schemas.enable的设置如果value.converter.schemas.enablefalseTopic 中的消息必须是纯 JSON例如{id:1}。如果value.converter.schemas.enabletrueTopic 中的消息必须包含schema和payload结构例如{schema:{...},payload:{id:1}}。若配置与实际消息格式不匹配Sink 会跳过或抛出转换异常。验证方式查看启动终端输出的日志中是否有DataException或ParseException。4.检查文件路径与写权限确认输出文件创建成功。检查connect-file-sink.properties中的file路径确认指定路径的上级目录是否存在且当前运行 Kafka 的用户拥有写权限。如果写的是相对路径它会创建在执行启动命令时的当前工作目录下。验证方式使用绝对路径测试例如file/home/apps/access.log.sink并检查该路径下是否成功创建了文件。5.检查 logs 目录下的日志文件定位后台日志。-daemon模式默认会将控制台标准输出和错误日志重定向到 Kafka 安装目录下的logs/文件夹中例如connect.out或connect.log。