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

使用CloudPubSubSourceConnector同步GCP PubSub消息到Kafka出现乱码问题

问题原因

  1. 配置中同时存在互斥的GCP凭证参数gcp.credentials.file.path和gcp.credentials.json,其中gcp.credentials.json参数需要传入完整的GCP服务账号JSON凭证字符串,你填写的是文件名gcp.json属于无效取值,导致连接器运行时偶发加载到错误凭证,消息处理进入异常分支。
  2. 异常分支返回的消息值为byte[]类型,你配置的org.apache.kafka.connect.storage.StringConverter仅支持String类型输入,遇到byte[]类型时会直接调用其默认toString()方法,生成[B@xxxx格式的无效字符串写入Kafka。
  3. 当连接器优先加载到gcp.credentials.file.path对应的正确凭证时,消息处理流程正常,返回String类型的消息值,写入Kafka的内容正常,因此出现一半概率正常一半概率异常的现象。

解决方案

第一步:修正凭证配置

删除配置中错误的"gcp.credentials.json":"gcp.json"行,仅保留gcp.credentials.file.path参数即可,避免凭证加载冲突。

第二步:修复类型转换问题,二选一即可

方案A:保留StringConverter,强制连接器返回RAW格式字符串

在连接器配置中新增以下两项:

"cps.message.format": "RAW",
"value.converter.encoding": "UTF-8"

该方案无需修改消费端代码,连接器会统一将Pub/Sub消息的payload转换为UTF-8编码的字符串写入Kafka。

方案B:替换为ByteArrayConverter,消费端自行解码

将配置中的value转换器替换为字节数组转换器:

"value.converter": "org.apache.kafka.connect.converters.ByteArrayConverter"

该方案稳定性更高,消费端读取到Kafka消息的字节数组后,自行按UTF-8编码解码为字符串即可。

内容的提问来源于stack exchange,提问作者paul

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 22:15:05