【发布时间】:2022-01-04 15:02:19
【问题描述】:
我正在尝试使用数据块读取 Azure 事件中心数据。
我有一个在 nodejs 中运行的生产者,以及一个用于测试的消费者(在不同的消费者组上),并且一切似乎都运行良好。
我在 databricks 中使用以下 pyspark 代码来获取数据:
# Initialize event hub config dictionary with connectionString
connectionString = "Endpoint=sb://XXX.servicebus.windows.net/;SharedAccessKeyName=test;SharedAccessKey=XXX;EntityPath=XXX"
ehConf = {}
ehConf['eventhubs.connectionString'] = sc._jvm.org.apache.spark.eventhubs.EventHubsUtils.encrypt(connectionString)
# Add consumer group to the ehConf dictionary
ehConf['eventhubs.consumerGroup'] = "databricks"
# Read events from the Event Hub
messages = spark.readStream.format("eventhubs").options(**ehConf).load()
# Visualize the Dataframe in realtime
display(messages)
问题在于,如果在笔记本运行时发送数据,它只会从流中读取数据。如果我生成数据然后运行笔记本,它不会出现。
我错过了什么?我想用它每隔一小时左右从流中收集数据并保存。
配置:
Databricks 运行时:7.3LTS(Spark 3.0.1、Scala 2.12)
Azure eventthub 库:com.microsoft.azure:azure-eventhubs-spark_2.12:2.3.21
【问题讨论】:
标签: pyspark azure-databricks azure-eventhub