PySpark 2.4消费Kafka中JSON数据时value转string乱码问题
问题根因与解决方案
问题结论
你的消费代码语法没有错误,但核心问题是错误假设Kafka中存储的value是纯UTF-8编码的JSON文本,直接将二进制格式的value强转为string必然出现乱码。
根因分析
- 从你贴出的value二进制样例看,所有消息都以
00 00 00 05开头,后续紧跟C0起始的字节,这部分不属于合法UTF-8文本编码范围,是二进制协议头,不属于JSON内容:- 如果PG到Kafka的同步链路用了Confluent Schema Registry + Avro/Protobuf序列化,前5字节是固定序列化头:1字节魔数
0x00+ 4字节整型的Schema ID(样例中值为5),后续内容是二进制序列化的结构化数据,本身就不是纯JSON文本,直接按UTF-8解码必然乱码。 - 如果同步任务直接将PostgreSQL原生逻辑复制(pgoutput)的二进制流写入Kafka,前面的字节是PG逻辑复制协议的固定元数据头(包含消息类型、事务ID、表映射关系等二进制信息),也不是纯文本JSON。
- 如果PG到Kafka的同步链路用了Confluent Schema Registry + Avro/Protobuf序列化,前5字节是固定序列化头:1字节魔数
- 你代码中
df['value'].cast('string')的逻辑,是将整个二进制字节数组按JVM默认UTF-8编码全量转为字符串,非文本的二进制头会被解码为不可识别的乱码字符,就会出现你看到的�����īVAL效果;后面能看到VAL、login这类正常字符,是因为这部分刚好是ASCII编码的文本内容,可以被UTF-8正常解码。 - 因为转成的字符串开头是乱码,不符合JSON语法,后续
from_json会直接返回null,最终na.drop会把所有数据全部过滤掉,无法得到正常解析结果。
消费侧解决方法
因为你无法修改生产者逻辑,可以按以下步骤处理:
- 先确认链路序列化配置
找同步链路的运维/开发负责人确认PG到Kafka的同步规则:- 如果是用Avro/Protobuf序列化,不要直接把value转string,使用Spark对应的
from_avro/from_protobuf函数,提前从Schema Registry拉取对应ID(样例中Schema ID为5)的Schema,直接对二进制value做反序列化提取结构化字段即可。 - 如果确认写入的是纯JSON,只是前面带固定长度的二进制冗余头,先裁剪二进制前缀再做解析。
- 如果是用Avro/Protobuf序列化,不要直接把value转string,使用Spark对应的
- 裁剪固定二进制头(纯JSON场景)
先对比你预期的JSON文本的十六进制编码,和实际拿到的value二进制内容,确定冗余头的固定长度,用Spark内置函数裁剪二进制内容后再转字符串。比如经确认头固定为10字节,可使用如下逻辑转换:from pyspark.sql import functions as F # 注意Spark中substring处理二进制类型时索引从1开始,裁掉前10字节就从第11位开始截取 df = df.select( F.expr("substring(value, 11)").cast("string").alias("value") ) # 后续再执行from_json解析即可 - 临时校验方法
可以先取单条消息的value二进制,手动裁掉不同长度的前缀后转字符串,直到得到合法的JSON文本,再确定固定裁剪长度即可。
内容的提问来源于stack exchange,提问作者amine jisung
相关产品推荐
相关产品推荐

