【问题标题】:Azure Databricks only gets Event Hub Data sent while its runnngAzure Databricks 仅在运行时获取发送的 Eventhub 数据
【发布时间】: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


    【解决方案1】:

    你有两个问题:

    1. display 默认使用临时检查点,所以当你下次运行它时,它不知道从哪里继续,所以它一次又一次地开始。如果您继续使用display,请将`checkpointLocation="some_path" 添加到显示调用中(参见docs)
    2. 默认情况下,EventHubs 连接器仅读取新数据(这就是它仅在运行时消耗数据的原因)-如果您想从头开始使用数据(仅在初始调用时)-您需要添加选项 eventhubs.startingPositions 编码开始职位(见doc) - 要从主题的开头开始阅读,请为此选项分配以下内容:
    import json
    
    startingEventPosition = {
      "offset": -1,  
      "seqNo": -1,            #not in use
      "enqueuedTime": None,   #not in use
      "isInclusive": True
    }
    
    ehConf["eventhubs.startingPosition"] = json.dumps(startingEventPosition)
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-09-17
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多