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

如何在PySpark中读取Azure Event Hub指定分区

读取Azure Event Hub指定分区的PySpark解决方案

要在PySpark中读取Azure Event Hub的指定分区(比如Partition 1),只需在Event Hub的配置参数中添加分区指定项即可,以下是修改后的代码及说明:

修改后的代码

# Config
connectionString = "Endpoint=sb://abcd"
eventHubName = "event_hub"

ehConf = {
  'eventhubs.connectionString' : sc._jvm.org.apache.spark.eventhubs.EventHubsUtils.encrypt(connectionString),
  'eventhubs.eventHubName': eventHubName,
  # 明确指定要读取的分区ID(Partition 1)
  'eventhubs.partition.id': "1",
  # 可选:配置该分区的起始消费位置,示例为从最新消息开始
  'eventhubs.startingPosition': '{"1": "-1"}'
}

df_stream  = spark.readStream.format("eventhubs")\
  .options(**ehConf)\
  .option("compression", "gzip") \
  .load() 

关键配置说明

  • eventhubs.partition.id:字符串类型,指定要消费的单个分区ID,这里设置为"1"对应目标Partition 1。添加该配置后,Spark流只会从指定分区拉取数据。
  • eventhubs.startingPosition:可选配置,以JSON字符串格式指定分区的起始消费偏移:
    • "-1":从分区最新的消息开始消费
    • "-2":从分区最早的消息开始消费
    • 也可以填入具体的偏移量数值(如"15678"),从指定位置开始消费

如果需要同时消费多个特定分区,可移除eventhubs.partition.id,在eventhubs.startingPosition中列出多个分区的配置(例如'{"1": "-1", "3": "-2"}'),同时调整Spark作业的并行度以匹配分区数量。

内容的提问来源于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 18:17:17