PySpark Streaming反序列化Kafka JSON消息:ID字段编码异常问题
解决PySpark消费Kafka时Decimal类型ID字段的解码问题
我来帮你搞定这个ID字段的解码问题!首先得搞清楚为什么ID会显示成AOo=这样的字符串:
你的消息是Kafka Connect通过JsonConverter生成的带Schema的JSON格式,其中ID字段被定义为org.apache.kafka.connect.data.Decimal(scale=0),这种类型会被序列化成base64编码的二进制字节,所以payload里的ID就是这些字节的base64字符串,我们需要把它转成对应的数值。
下面给你两种解决方案,一种适配你当前用的DStream API,另一种是更推荐的Structured Streaming API:
方法1:适配现有DStream代码,自定义valueDecoder
我们可以写一个自定义的解码函数,传给createDirectStream的valueDecoder参数,完成从base64字符串到整数的转换:
import base64 import json from pyspark.sql import SparkSession from pyspark.streaming import StreamingContext from pyspark.streaming.kafka import KafkaUtils def decode_kafka_value(value_bytes): # 把Kafka消息的字节解析成JSON对象 message = json.loads(value_bytes.decode('utf-8')) payload = message['payload'] # 解码ID的base64字符串为字节数组,再转成整数(scale=0等价于整数) # 这里用大端序(big-endian)是因为Kafka Connect的Decimal默认用大端存储 id_bytes = base64.b64decode(payload['ID']) payload['ID'] = int.from_bytes(id_bytes, byteorder='big', signed=False) # 返回schema和处理后的payload,你也可以只返回payload return message['schema'], payload if __name__ == "__main__": spark = SparkSession.builder.appName("my app").getOrCreate() sc = spark.sparkContext sc.setLogLevel('WARN') ssc = StreamingContext(sc, 5) kafka_params = { "bootstrap.servers": "kafkahost:9092", "group.id": "Deserialize" } # 传入自定义的valueDecoder kafka_stream = KafkaUtils.createDirectStream( ssc, ['mytopic'], kafka_params, valueDecoder=decode_kafka_value ) # 现在打印的payload里ID就是整数了 kafka_stream.foreachRDD(lambda rdd: rdd.foreach(lambda x: print(x[1]))) ssc.start() ssc.awaitTermination()
方法2:推荐使用Spark Structured Streaming(更现代的API)
DStream是Spark旧的流处理API,Structured Streaming提供了更简洁的DataFrame/SQL接口,也更容易维护。这里是实现代码:
from pyspark.sql import SparkSession from pyspark.sql.functions import from_json, col, udf from pyspark.sql.types import StructType, StructField, StringType, IntegerType import base64 # 自定义UDF,把base64编码的Decimal(scale=0)转成整数 def decode_decimal_base64(base64_str): if not base64_str: return None byte_data = base64.b64decode(base64_str) return int.from_bytes(byte_data, byteorder='big', signed=False) decode_decimal_udf = udf(decode_decimal_base64, IntegerType()) if __name__ == "__main__": spark = SparkSession.builder.appName("KafkaDecimalDecode").getOrCreate() spark.sparkContext.setLogLevel('WARN') # 读取Kafka流数据 kafka_df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "kafkahost:9092") \ .option("subscribe", "mytopic") \ .option("group.id", "Deserialize") \ .load() # 定义Kafka消息的Schema结构 message_schema = StructType([ StructField("schema", StringType()), StructField("payload", StructType([ StructField("ID", StringType()), # 先按字符串读取base64内容 StructField("COLUMN1", StringType()) ])) ]) # 解析JSON格式的消息内容 parsed_df = kafka_df.select( from_json(col("value").cast(StringType()), message_schema).alias("message") ).select("message.payload.*") # 把ID字段解码成整数 final_df = parsed_df.withColumn("ID", decode_decimal_udf(col("ID"))) # 输出到控制台查看结果(生产环境可以换成其他输出源) query = final_df.writeStream \ .outputMode("append") \ .format("console") \ .start() query.awaitTermination()
额外说明
- 如果你遇到的Decimal字段
scale不是0(比如scale=2代表两位小数),可以修改解码逻辑:导入decimal库,把整数转成Decimal后除以10^scale,示例代码:import decimal def decode_decimal_base64(base64_str, scale=2): if not base64_str: return None byte_data = base64.b64decode(base64_str) return decimal.Decimal(int.from_bytes(byte_data, byteorder='big', signed=False)) / (10 ** scale) - 你提供的示例中
AOo=解码后是字节b'\x00\x0e',转成整数就是14,这就是你原本的ID数值。
内容的提问来源于stack exchange,提问作者anonuser1234
相关产品推荐
相关产品推荐

