You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Azure Databricks与Azure Event Hubs的Python连接及键值对参数发送问题求助

Fixing Azure Databricks + Event Hubs Python Integration & Sending Key-Value Events

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.11 packages—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
  • Keep org.apache.spark:spark-avro_2.12:3.1.1 only 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 Send permissions on the Event Hub (use a Shared Access Policy with Send rights).
  • 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.27 17:37:40