Kafka至Snowflake流同步失败:JSON解析错误求助
- Kafka 3.9.0 / Snowflake Kafka连接器 3.1.2
- 也曾尝试 Kafka 3.3.1 / Snowflake Kafka连接器 2.0.0
- Snowflake Cloud
- 需实现Java应用 → Kafka → Snowflake的流式数据传输
- Java应用到Kafka的数据流正常,可在Topic中查看数据
Java应用发送到Kafka的示例数据(Topic中可见):
{"SalesDocumentItemCategory":"ASD","BillingRelevanceCode":"CODE1","ScheduleLineIsAllowed":"Y","PricingRelevance":"NA","TextDeterminationProcedure":"99","PartnerDeterminationProcedure":"N","PropagatePrftbltySgmt2BOM":"BOM1","CostDeterminationIsRequired":"X"} {"SalesDocumentItemCategory":"DAS","BillingRelevanceCode":"CODE2","ScheduleLineIsAllowed":"Y","PricingRelevance":"NA","TextDeterminationProcedure":"99","PartnerDeterminationProcedure":"N","PropagatePrftbltySgmt2BOM":"BOM1","CostDeterminationIsRequired":"X"}
- 从Kafka到Snowflake使用Snowflake Kafka连接器结合Snowpipe Streaming
- 运行连接器时出现错误(见下文)
- 因需要自动schema化,采用Snowpipe Streaming导入方式
- 连接器报错后,Snowflake中已创建表,但无行和列
- 用相同数据通过kafka-console-producer手动生产时,数据能按预期schema写入Snowflake
是否遗漏了连接器的某些配置?
连接器属性配置
name=kfka_sf_connector3 topics=TP_BILLINGDOCUMENT connector.class=com.snowflake.kafka.connector.SnowflakeSinkConnector tasks.max=8 buffer.count.records=10000 buffer.flush.time=60 buffer.size.bytes=5000000 snowflake.url.name=xxx.snowflakecomputing.com snowflake.user.name=xxx snowflake.private.key=<XXXX> snowflake.database.name=DB_TRFT snowflake.schema.name=SCH_SAL1 #key.converter=org.apache.kafka.connect.json.JsonConverter key.converter=com.snowflake.kafka.connector.records.SnowflakeJsonConverter #key.converter=org.apache.kafka.connect.storage.StringConverter #key.converter=org.apache.kafka.connect.converters.ByteArrayConverter key.converter.schemas.enable=false #value.converter=org.apache.kafka.connect.json.JsonConverter value.converter=com.snowflake.kafka.connector.records.SnowflakeJsonConverter #value.converter=org.apache.kafka.connect.storage.StringConverter #value.converter=org.apache.kafka.connect.converters.ByteArrayConverter value.converter.schemas.enable=false snowflake.ingestion.method=SNOWPIPE_STREAMING snowflake.enable.schematization=TRUE snowflake.role.name=kafka_connector_role
错误信息
ERROR [dsp_kfka_sf_connector3|task-7] WorkerSinkTask{id=kfka_sf_connector3-7} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted (org.apache.kafka.connect.runtime.WorkerTask:234) org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:261) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:180) at org.apache.kafka.connect.runtime.WorkerSinkTask.convertAndTransformRecord(WorkerSinkTask.java:538) at org.apache.kafka.connect.runtime.WorkerSinkTask.convertMessages(WorkerSinkTask.java:518) at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:344) at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:247) at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:216) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:226) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:281) at org.apache.kafka.connect.runtime.isolation.Plugins.lambda$withClassLoader$1(Plugins.java:238) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:572) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:317) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1144) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:642) at java.base/java.lang.Thread.run(Thread.java:1583) Caused by: org.apache.kafka.connect.errors.DataException: Converting byte[] to Kafka Connect data failed due to serialization error: at org.apache.kafka.connect.json.JsonConverter.toConnectData(JsonConverter.java:347) at org.apache.kafka.connect.storage.Converter.toConnectData(Converter.java:91) at org.apache.kafka.connect.runtime.WorkerSinkTask.lambda$convertAndTransformRecord$3(WorkerSinkTask.java:538) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:208) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:245) Caused by: org.apache.kafka.common.errors.SerializationException: com.fasterxml.jackson.core.JsonParseException: Unrecognized token 'AES': was expecting (JSON String, Number, Array, Object or token 'null', 'true' or 'false') at [Source: REDACTED (`StreamReadFeature.INCLUDE_SOURCE_IN_LOCATION` disabled); line: 1, column: 4] at org.apache.kafka.connect.json.JsonDeserializer.deserialize(JsonDeserializer.java:75) at org.apache.kafka.connect.json.JsonConverter.toConnectData(JsonConverter.java:345) Caused by: com.fasterxml.jackson.core.JsonParseException: Unrecognized token 'ASD': was expecting (JSON String, Number, Array, Object or token 'null', 'true' or 'false') at [Source: REDACTED (`StreamReadFeature.INCLUDE_SOURCE_IN_LOCATION` disabled); line: 1, column: 4] at com.fasterxml.jackson.core.JsonParser._constructError(JsonParser.java:2481) at com.fasterxml.jackson.core.base.ParserMinimalBase._reportError(ParserMinimalBase.java:762) at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._reportInvalidToken(UTF8StreamJsonParser.java:3703) at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._handleUnexpectedValue(UTF8StreamJsonParser.java:2791) at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._nextTokenNotInObject(UTF8StreamJsonParser.java:911) at com.fasterxml.jackson.core.json.UTF8StreamJsonParser.nextToken(UTF8StreamJsonParser.java:797) at com.fasterxml.jackson.databind.ObjectMapper._readTreeAndClose(ObjectMapper.java:4928) at com.fasterxml.jackson.databind.ObjectMapper.readTree(ObjectMapper.java:3292) at org.apache.kafka.connect.json.JsonDeserializer.deserialize(JsonDeserializer.java:73)
核心问题分析
错误日志显示JSON解析失败,提示Unrecognized token 'ASD',说明连接器接收到的消息格式不符合预期。对比手动生产数据成功、Java应用生产失败的情况,问题出在Java应用发送的Kafka消息格式:
示例数据中,Java应用把两个JSON对象直接拼接在同一条消息里(无分隔符、非JSON数组),而Kafka Connect的JSON转换器(包括Snowflake的转换器)默认要求每条消息是单个合法JSON对象/数组。手动生产时,每个JSON作为单独消息发送,因此能成功。
另外,错误栈显示连接器仍在使用org.apache.kafka.connect.json.JsonConverter,可能是配置加载异常,或SnowflakeJsonConverter依赖底层JSON解析逻辑导致格式校验严格。
修复步骤
1. 修正Java应用的消息发送格式
确保Java应用发送到Kafka的每条消息是单个合法JSON对象,不要将多个JSON拼接成一条消息。若需批量发送,使用JSON数组格式:
[{"SalesDocumentItemCategory":"ASD", ...}, {"SalesDocumentItemCategory":"DAS", ...}]
或直接拆分每个JSON为独立的Kafka消息。
2. 调整连接器转换器配置
简化配置,移除注释行避免干扰,明确使用合适的转换器:
name=kfka_sf_connector3 topics=TP_BILLINGDOCUMENT connector.class=com.snowflake.kafka.connector.SnowflakeSinkConnector tasks.max=8 buffer.count.records=10000 buffer.flush.time=60 buffer.size.bytes=5000000 snowflake.url.name=xxx.snowflakecomputing.com snowflake.user.name=xxx snowflake.private.key=<XXXX> snowflake.database.name=DB_TRFT snowflake.schema.name=SCH_SAL1 key.converter=org.apache.kafka.connect.storage.StringConverter key.converter.schemas.enable=false value.converter=com.snowflake.kafka.connector.records.SnowflakeJsonConverter value.converter.schemas.enable=false snowflake.ingestion.method=SNOWPIPE_STREAMING snowflake.enable.schematization=TRUE snowflake.role.name=kafka_connector_role
消息key无需处理时,使用StringConverter更简洁;value用SnowflakeJsonConverter适配Snowflake的schema化需求。
3. 验证消息格式
用kafka-console-consumer查看消息实际格式,确认每条消息为单个合法JSON:
kafka-console-consumer.sh --bootstrap-server <你的Kafka Broker地址> --topic TP_BILLINGDOCUMENT --from-beginning
若输出每条消息都是单独的JSON对象,说明格式正确;若为多个JSON拼接,需调整Java应用的生产者逻辑。
4. 检查Snowpipe Streaming权限
确保kafka_connector_role拥有足够权限:
- 数据库、schema的USAGE权限
- 表的CREATE、INSERT权限
- Snowpipe Streaming相关权限(如
STREAMING_INSERT)
内容的提问来源于stack exchange,提问作者Nikhil Nair

