【问题标题】:Is Spark Streaming availableNow trigger compatible with Azure Event Hub?Spark Streaming availableNow 触发器是否与 Azure 事件中心兼容?
【发布时间】: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 触发器是否损坏,或者我在这里做错了什么?

【问题讨论】:

标签: apache-spark pyspark databricks azure-eventhub


【解决方案1】:

'avaiableNow' 触发器似乎尚未在'azure-event-hub-spark' 包中实现。

但是有一个解决方法可以使用 Kafka 连接器连接到 Azure 事件中心 - https://github.com/Azure/azure-event-hubs-for-kafka/tree/master/tutorials/spark

所以基本上前面的代码变成了

bootstrap_servers = "my-evh-namespace.servicebus.windows.net:9093"
eventhub_endpoint = "my-evh-namespace-endpoint"

# The 'kafkashaded' part here is because it's running in Databricks.
# Otherwise drop that part.
EH_SASL = f"kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModule required username="$ConnectionString" password="{eventhub_endpoint}";"

topic = "my-eventhub-name"

read_stream = spark.readStream 
    .format("kafka") 
    .option("kafka.bootstrap.servers", bootstrap_servers) 
    .option("kafka.sasl.mechanism", "PLAIN") 
    .option("kafka.security.protocol", "SASL_SSL") 
    .option("kafka.sasl.jaas.config", EH_SASL) 
    .option("subscribe", topic) 
    .option("maxOffsetsPerTrigger", 1000) 
    .option("startingOffsets", "earliest") 
    .option("includeHeaders", "true") 
    .load()

# Notice that the output writeStream remains the same.
stream = read_stream.writeStream 
  .format("delta") 
  .option("checkpointLocation", checkpoint_location) 
  .trigger(availableNow=True) 
  .toTable(full_table_name, mode="append")

这会导致流按预期执行 - 以 maxOffsetsPerTrigger 大小的批次摄取直到开始时间之前的所有事件

【讨论】:

    猜你喜欢
    • 2016-11-13
    • 2017-10-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-03-19
    • 2023-03-25
    • 2021-11-27
    • 1970-01-01
    相关资源
    最近更新 更多