Azure Databricks与Azure Event Hubs的Python连接及键值对参数发送问题求助
First, let's clear up the root cause of your unable to import from Azure Event hubs error: the libraries you've installed are Spark/Scala connectors for Event Hubs, but if you're trying to use the native Azure Event Hubs Python SDK, you're missing the required PyPI package. Alternatively, if you want to mirror your Scala implementation using Spark Structured Streaming, the import approach is different than the raw Python SDK.
Let's break down both valid approaches, along with fixes for your library setup:
1. Use Spark Structured Streaming (Mirror Your Scala Workflow)
This is the closest match to your successful Scala implementation, and uses the Spark connector you've already started installing.
Correct Library Setup
- Remove Scala 2.11 libraries: Databricks Runtime 7.3 LTS and newer use Scala 2.12, so scrap any
azure-eventhubs-spark_2.11packages—they'll cause version conflicts. - Match connector version to your Databricks Runtime:
- For Runtime 11.3 LTS (Spark 3.3.0, Scala 2.12): Use
com.microsoft.azure:azure-eventhubs-spark_2.12:2.3.22 - For Runtime 10.4 LTS (Spark 3.2.1, Scala 2.12): Use
com.microsoft.azure:azure-eventhubs-spark_2.12:2.3.19
- For Runtime 11.3 LTS (Spark 3.3.0, Scala 2.12): Use
- Keep
org.apache.spark:spark-avro_2.12:3.1.1only if you're working with Avro data; it's not required for basic key-value events.
Python Code to Send Key-Value Events
You don't need to import azure.eventhub here—we use Spark's DataFrame API to write directly to Event Hubs:
# Configure Event Hub credentials eh_connection_string = "YOUR_EVENT_HUB_CONNECTION_STRING" eh_entity_path = "YOUR_EVENT_HUB_NAME" # Encrypt the connection string (required for Spark connector) eh_config = { "eventhubs.connectionString": sc._jvm.org.apache.spark.eventhubs.EventHubsUtils.encrypt(eh_connection_string), "eventhubs.entityPath": eh_entity_path } # Create a DataFrame with your key-value pair data key_value_data = [ {"user_id": "123", "action": "login", "timestamp": "2024-05-20T14:30:00"}, {"user_id": "456", "action": "purchase", "timestamp": "2024-05-20T14:31:00"} ] df = spark.createDataFrame(key_value_data) # Write the DataFrame to Event Hubs df.write.format("eventhubs").options(**eh_config).mode("append").save()
2. Use Native Azure Event Hubs Python SDK
If you prefer to use the raw Python SDK (not Spark), you'll need to install the PyPI package alongside your Spark libraries.
Install Required PyPI Package
Run this in a Databricks notebook cell to install the SDK:
%pip install azure-eventhub==5.11.0
(Pick a version compatible with your Databricks Runtime—5.11.0 works well with most recent LTS runtimes.)
Python Code to Send Key-Value Events
Now you can import the SDK and send structured JSON events:
from azure.eventhub import EventHubProducerClient, EventData import json # Event Hub credentials connection_str = "YOUR_EVENT_HUB_CONNECTION_STRING" eventhub_name = "YOUR_EVENT_HUB_NAME" # Initialize producer client producer = EventHubProducerClient.from_connection_string(connection_str, eventhub_name=eventhub_name) # Prepare key-value event (convert to JSON string for structured delivery) with producer: event_batch = producer.create_batch() # Add multiple key-value events to the batch event_batch.add(EventData(json.dumps({"product": "laptop", "price": 999.99, "in_stock": True}))) event_batch.add(EventData(json.dumps({"product": "mouse", "price": 29.99, "in_stock": False}))) # Send the batch producer.send_batch(event_batch)
Critical Troubleshooting Tips
- Avoid mixed library versions: Never mix Scala 2.11 and 2.12 packages in the same cluster—this causes unresolved dependencies.
- Connection string permissions: Ensure your connection string has
Sendpermissions on the Event Hub (use a Shared Access Policy withSendrights). - Splunk compatibility: Both approaches send structured data (either Spark's serialized rows or JSON strings), which Splunk can parse by configuring a JSON sourcetype for your Event Hub input.
内容的提问来源于stack exchange,提问作者Vignesh Ravichandran

