如何在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
相关产品推荐
相关产品推荐

