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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 04:32:25