You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.13 09:39:50