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

如何在Fabric中用PySpark读取Azure Event Hub中的压缩数据

解决PySpark读取Event Hub手动压缩+Base64编码数据返回Null的问题

问题根源

你手动对数据执行了gzip压缩→Base64编码的两步处理,但PySpark Event Hub连接器的compression选项仅适用于Event Hub服务端自动压缩的场景,无法识别客户端手动编码的数据,因此直接配置该选项无效,导致解析后body为Null。

正确处理流程

需要在PySpark中反向执行编码步骤:先对body列做Base64解码,再解gzip压缩,最后解析JSON。

修改后的完整读取代码

  1. 读取Event Hub流数据(无需指定compression选项):
df_stream = spark.readStream.format("eventhubs")\
  .options(**ehConf)\
  .load() 
  1. 使用Spark内置函数完成解码+解压缩(比UDF更高效):
import pyspark.sql.functions as F
from pyspark.sql.types import StringType

# 定义解码+解压缩逻辑
def decompress_gzip_base64(col):
    # 1. 将body字符串解码为二进制数据
    decoded_bytes = F.unbase64(col)
    # 2. 对二进制数据执行gzip解压缩,再转为字符串
    return F.uncompress(decoded_bytes, "gzip").cast(StringType())

# 处理原始body列,得到可解析的JSON字符串
df_stream_processed = df_stream.withColumn(
    "decoded_body",
    decompress_gzip_base64(F.col("body").cast("string"))
)

# 解析JSON为结构化数据
df_stream_body = df_stream_processed.select(
    F.from_json(F.col("decoded_body"), message_schema).alias("Payload")
)
  1. 更新foreach_batch_function逻辑:
def foreach_batch_function(df_stream, epoch_id):
    # 先完成解码和解压缩
    df_stream_processed = df_stream.withColumn(
        "decoded_body",
        decompress_gzip_base64(F.col("body").cast("string"))
    )
    # 解析JSON并创建临时视图
    df_stream_body = df_stream_processed.select(
        F.from_json(F.col("decoded_body"), message_schema).alias("Payload")
    )
    df_stream_body.createOrReplaceTempView("stream_temp")
    # 后续更新Delta表的逻辑...

关键注意事项

  • 严格遵循Base64解码→Gzip解压缩→JSON解析的顺序,不能颠倒。
  • 若遇到Spark内置函数兼容性问题,可改用UDF实现(需额外引入base64、gzip、BytesIO库):
from io import BytesIO
import gzip
import base64

@F.udf(StringType())
def decompress_gzip_base64_udf(body_str):
    if not body_str:
        return None
    decoded_bytes = base64.b64decode(body_str)
    with gzip.GzipFile(fileobj=BytesIO(decoded_bytes), mode='rb') as f:
        return f.read().decode('utf-8')

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 23:42:38