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

无法在Databricks查看Event Hub生成数据:字节数组异常及后续问题

解决Databricks读取Event Hubs模拟Yellow Taxi数据流的错误

第一个错误:Input byte array has wrong 4-byte ending unit

这个错误的核心原因是Event Hubs的Generate Data功能生成的消息体是Base64编码的二进制数据,直接通过cast("string")转换会触发Base64解码格式错误。正确的处理步骤是先对二进制body做Base64解码,再转为字符串。

修改后的代码如下:

from pyspark.sql.functions import unbase64, col

connectionString = 'Endpoint=sb://carlsspace.servicebus.windows.net/;SharedAccessKeyName=myaccesspolicy;SharedAccessKey=lfZTXXXXXXXXXXXXXXXXXXXVaNv9wRM+AEhJ3XsFc=;EntityPath=carlshub'
ehConf = {}

# 适配2.3.15+版本的加密要求
ehConf['eventhubs.connectionString'] = sc._jvm.org.apache.spark.eventhubs.EventHubsUtils.encrypt(connectionString)

df = spark \
  .readStream \
  .format("eventhubs") \
  .options(**ehConf) \
  .load()

# 先解码Base64,再转换为字符串
df = df.withColumn("decoded_body", unbase64(col("body")).cast("string"))
display(df.select("decoded_body"))

连接字符串加密的关键注意事项

  • 严格匹配Spark Event Hubs Connector版本:2.3.15及以上版本必须使用EventHubsUtils.encrypt()加密连接字符串;低于该版本直接传入明文连接串即可。
  • Databricks社区版默认Spark版本通常较高,务必遵循高版本的加密要求,否则会触发授权或格式类错误。

额外排查方向

  • 验证共享访问权限:确保myaccesspolicy拥有Event Hub的Listen权限,无权限会导致连接失败。
  • 核对EntityPath:必须与你的Event Hub名称完全一致,注意大小写敏感。
  • 查看流监控日志:在Databricks中点击display生成的流图表,进入监控界面查看完整错误栈,这能帮你定位后续出现的新问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 08:33:28