Kafka Sink Connector执行TASK_PUT阶段报错,数据无法写入数据库
Kafka Protobuf Sink Connector TASK_PUT报错无数据写入排查方案
问题场景
开发的Kafka Sink Connector用于将Protobuf主题数据推送至SQL Server数据库,连接器状态显示运行中,但日志显示TASK_PUT阶段报错,无数据写入目标表RealTimeEvents.connector.SEG_32RED_ALL_DEPOSITS_PROTO。
关键信息
报错日志
[2023-03-21 22:26:47,367] ERROR Error encountered in task 32red-all-deposit-proto-to-mssql-flat-test-postgres-3. Executing stage 'TASK_PUT' with class 'org.apache.kafka.connect.sink.SinkTask', where consumed record is {topic='STR_32RED_ALL_DEPOSITS_FLAT', partition=3, offset=116, timestamp=1679399268996, timestampType=CreateTime}. (org.apache.kafka.connect.runtime.errors.LogReporter)
Protobuf主题样本数据(JSON格式)
{ "EVENTSEQUENCENUMBER": "31851094", "GAMINGSYSTEMID": 323, "ROUTERID": 3, "USERNAME": "nonapplicable", "USERID": "1516780", "PRODUCTID": 380, "SESSIONPRODUCTID": 380, "SESSIONID": 190985932, "CURRENCYISOCODE": "GBP", "OPERATORCURRENCYISOCODE": "GBP", "PLAYERTOOPERATOREXCHANGERATE": 1.0, "DEPOSITAMOUNT": 10.0, "DEPOSITTYPE": "Ecash", "DEPOSITMETHOD": "Pay_Pal", "BALANCEAFTERDEPOSIT": 10.16, "ECASHPROCESSORID": 160, "ECASHPROCESSORDESCRIPTION": "Pay_Pal", "ISSUCCESS": 1, "UTCEVENTTIME": "2023-03-21 11:47:25.5677383", "TICKSEVENTTIME": "638149960455677383", "COUNTRYLONGCODE": "GBR", "LANGUAGECODE": "EN", "SESSIONCOUNTRYLONGCODE": "GBR", "NUMDEPOSITSTOTAL": 2036, "TOTALDEPOSITSAMOUNT": 20901.0, "TRANSACTIONID": 11369506, "TRANSACTIONUTCDATETIME": "2023-03-21 11:47:25.510", "TRANSACTIONNUMBER": "380_11369506", "TRANSACTIONSTATUSID": 1, "TRANSACTIONSTATUS": "Accept", "EVENTID": 6026, "EVENTNAME": "Pay_Pal_Transaction - Automatic Update", "TRANSACTIONEVENTTYPEID": 1, "TRANSACTIONEVENTTYPE": "Purchase", "ECASHREQUESTTIME": "2023-03-21 11:47:08.303", "PLAYERGROUP": "7cf2209c-1500-495b-9e80-ffc2ae4a9902" }
Sink Connector配置
{ "name": "32red-all-deposit-proto-to-mssql-flat-test", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "table.name.format": "RealTimeEvents.connector.SEG_32RED_ALL_DEPOSITS_PROTO", "connection.password": "password", "topics": "STR_32RED_ALL_DEPOSITS_FLAT", "tasks.max": "4", "batch.size": "100", "max.retries": "3", "auto.create": "true", "auto.evolve": "true", "errors.tolerance": "all", "errors.log.enable": true, "connection.user": "RealTimeEvents", "name": "32red-all-deposit-proto-to-mssql-flat-test", "connection.url": "jdbc:sqlserver://BEL1DATAOPSDV01.data32red.com:1433;Catalog=RealTimeEvents;databaseName=RealTimeEvents", "insert.mode": "upsert", "pk.mode" : "none", "pk.fields" : "none", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "io.confluent.connect.protobuf.ProtobufConverter", "value.converter.schemas.enable": "true", "value.converter.schema.registry.url": "https://cp-schema-registry.d2-dev.aws.kindredgroup.com", "value.converter.schema.registry.ssl.truststore.location": "/etc/kafka-connect/secrets/cp-kafka-connect.d2-dev.aws.kindredgroup.com.truststore.jks", "value.converter.schema.registry.ssl.truststore.password": "${dir:/etc/kafka-connect/secrets:cp-kafka-connect.d2-dev.aws.kindredgroup.com.TrustStorePass}", "value.converter.schema.registry.ssl.keystore.location": "/etc/kafka-connect/secrets/cp-kafka-connect.d2-dev.aws.kindredgroup.com.keystore.jks", "value.converter.schema.registry.ssl.keystore.password": "${dir:/etc/kafka-connect/secrets:cp-kafka-connect.d2-dev.aws.kindredgroup.com.KeyStorePass}", "value.converter.schema.registry.ssl.key.password": "${dir:/etc/kafka-connect/secrets:cp-kafka-connect.d2-dev.aws.kindredgroup.com.KeyPass}", "errors.log.include.messages": true } }
排查与解决方案
1. 补全详细错误日志
当前日志仅提示TASK_PUT阶段报错,未给出具体异常信息。需调整日志级别获取完整堆栈:
- 修改Kafka Connect日志配置,将
org.apache.kafka.connect.runtime.errors和io.confluent.connect.jdbc的日志级别设为DEBUG或TRACE - 重启连接器后重新触发任务,查看完整异常堆栈,定位具体错误(如字段类型不匹配、数据库权限不足、Schema Registry解析失败等)
2. 修复Upsert模式配置冲突
配置中insert.mode: upsert但pk.mode: none存在矛盾:
- Upsert模式必须指定主键(用于判断插入/更新操作),否则无法生成合法的Upsert语句
- 解决方式:
- 若只需插入数据,将
insert.mode改为insert - 若需Upsert,指定合理主键,比如
pk.mode: record_value,pk.fields: TRANSACTIONID(样本数据中TRANSACTIONID为唯一标识)
- 若只需插入数据,将
3. 验证字段类型兼容性
自动建表时,Connector会根据Protobuf Schema推断数据库字段类型,可能存在不兼容场景:
- 检查
TICKSEVENTTIME:样本数据为字符串类型大数值,SQL Server可能推断为VARCHAR,实际应使用BIGINT,需手动调整表结构或确认auto.evolve是否自动适配(注意auto.evolve不支持所有类型变更) - 检查
UTCEVENTTIME等时间字段:确认Protobuf时间类型能否被正确解析为SQL Server的DATETIME2或TIMESTAMP类型,解析失败会导致插入报错
4. 确认Schema Registry连接有效性
Protobuf Converter依赖Schema Registry获取消息Schema,需确保:
- Connector节点能正常访问Schema Registry URL,可通过curl测试连通性
- SSL证书配置正确:验证
truststore、keystore文件路径及密码有效性,确认证书未过期 - 确认Schema Registry中存在该主题对应的Protobuf Schema,且版本正确
5. 检查数据库权限与表状态
- 确认
RealTimeEvents用户对RealTimeEvents.connectorschema拥有建表、插入/更新数据的权限 - 检查目标表
SEG_32RED_ALL_DEPOSITS_PROTO是否已自动创建,若未创建,排查auto.create配置是否因权限不足或Schema解析失败未生效
6. 临时调整错误容忍策略定位问题
当前errors.tolerance: all会跳过错误记录,无法定位根因:
- 临时改为
errors.tolerance: none,让连接器遇错即停,便于捕获完整错误信息 - 问题解决后再恢复原有错误容忍配置
内容的提问来源于stack exchange,提问作者mayanksweden
相关产品推荐
相关产品推荐

