【发布时间】:2022-10-13 22:38:15
【问题描述】:
我正在尝试使用带有availableNow 触发器的火花流将数据从 Azure 事件中心摄取到 Databricks 中的 Delta Lake 表中。
我的代码:
conn_str = "my conn string"
ehConf = {
"eventhubs.connectionString":
spark.sparkContext._jvm.org.apache.spark.eventhubs.EventHubsUtils.encrypt(conn_str),
"eventhubs.consumerGroup":
"my-consumer-grp",
}
read_stream = spark.readStream \
.format("eventhubs") \
.options(**ehConf) \
.load()
stream = read_stream.writeStream \
.format("delta") \
.option("checkpointLocation", checkpoint_location) \
.trigger(availableNow=True) \
.toTable(full_table_name, mode="append")
根据文档https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html#triggers
availableNow 触发器应该以微批处理方式处理当前可用的所有数据。
但是,这并没有发生,相反,它只处理 1000 行。 流的输出讲述了这个故事:
{
"sources" : [ {
"description" : "org.apache.spark.sql.eventhubs.EventHubsSource@2c5bba32",
"startOffset" : {
"my-hub-name" : {
"0" : 114198857
}
},
"endOffset" : {
"my-hub-name" : {
"0" : 119649573
}
},
"latestOffset" : {
"my-hub-name" : {
"0" : 119650573
}
},
"numInputRows" : 1000,
"inputRowsPerSecond" : 0.0,
"processedRowsPerSecond" : 36.1755236407047
} ]
}
我们可以清楚地看到偏移量的变化方式超过了 1000 次处理。
我已经验证了目标表的内容,它包含最后 1000 个偏移量。 \
根据 Pyspark https://github.com/Azure/azure-event-hubs-spark/blob/master/docs/PySpark/structured-streaming-pyspark.md#event-hubs-configuration 的事件中心配置maxEventsPerTrigger 默认设置为 1000*partitionCount,但这只会影响每批处理的事件数,而不影响 availableNow 触发器处理的记录总数。
使用触发器 once=True 运行相同的查询将改为摄取全部的事件(假设批量大小设置得足够大)。
Azure 事件中心的 availableNow 触发器是否损坏,或者我在这里做错了什么?
【问题讨论】:
-
我在 azure-event-hubs-spark github 上提出了一个关于此的问题。 github.com/Azure/azure-event-hubs-spark/issues/656 我怀疑他们还没有实现这个触发器支持。
标签: apache-spark pyspark databricks azure-eventhub