使用CloudPubSubSourceConnector同步GCP PubSub消息到Kafka出现乱码问题
问题原因
- 配置中同时存在互斥的GCP凭证参数
gcp.credentials.file.path和gcp.credentials.json,其中gcp.credentials.json参数需要传入完整的GCP服务账号JSON凭证字符串,你填写的是文件名gcp.json属于无效取值,导致连接器运行时偶发加载到错误凭证,消息处理进入异常分支。 - 异常分支返回的消息值为
byte[]类型,你配置的org.apache.kafka.connect.storage.StringConverter仅支持String类型输入,遇到byte[]类型时会直接调用其默认toString()方法,生成[B@xxxx格式的无效字符串写入Kafka。 - 当连接器优先加载到
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
相关产品推荐
相关产品推荐

