【问题标题】:Sending XML file content to Event Hub and read it from Databricks将 XML 文件内容发送到事件中心并从 Databricks 中读取
【发布时间】:2021-01-20 08:16:04
【问题描述】:

我正在尝试将 xml 文件(小于 100 kb)发送到 Azure 事件中心,然后在发送后读取 Databricks 中的事件。

现在我已经使用 Python SDK 以字节为单位发送 XML 的内容(这一步 WORKS)。但我想要实现的下一步是从事件的“主体”中读取 XML 内容,并使用 PYSPARK 创建一个 Spark Dataframe。

能够做到这一点,我有两个疑问:

1- 有没有我在spark.readStream 选项中指定事件“正文”的内容是 XML 的选项?

2- 是否有任何替代方法可以将该内容直接转储到 Spark Dataframe?

3- 将 XML 作为事件发送时缺少一些配置?

我正在尝试如下示例:

Python 事件制作者

# this is the python event hub message producer
import asyncio
from azure.eventhub.aio import EventHubProducerClient
from azure.eventhub import EventData
import xml.etree.ElementTree as ET
from lxml import etree
from pathlib import Path

connection_str= "Endpoint_str"
eventhub_name = "eventhub_name"

xml_path = Path("path/to/xmlfile.xml")

xml_data = ET.parse(xml_path)
tree = xml_data.getroot()
data = ET.tostring(tree)

async def run():
    # Create a producer client to send messages to the event hub.
    # Specify a connection string to your event hubs namespace and
    # the event hub name.
    producer = EventHubProducerClient.from_connection_string(conn_str=connection_str, eventhub_name=eventhub_name)
    async with producer:
        # Create a batch.
        event_data_batch = await producer.create_batch()

        # Add events to the batch.
        event_data_batch.add(EventData(data))

        # Send the batch of events to the event hub.
        await producer.send_batch(event_data_batch)

loop = asyncio.get_event_loop()
loop.run_until_complete(run())

事件阅读器

stream_data = spark \
    .readStream \
    .format('eventhubs') \
    .options(**event_hub_conf) \
    .option('multiLine', True) \
    .option('mode', 'PERMISSIVE') \
    .load()

谢谢!!!

【问题讨论】:

    标签: python xml azure apache-spark azure-eventhub


    【解决方案1】:

    所以我终于有了从事件中心主体读取 XML 的下一种方法。

    首先我使用import xml.etree.ElementTree as ET 库来解析XML 结构。

    stream_data = spark \
        .readStream \
        .format('eventhubs') \
        .options(**event_hub_conf) \
        .option('multiLine', True) \
        .option('mode', 'PERMISSIVE') \
        .load() \
        .select("body")
    
    df = stream_data.withColumn("body", stream_data["body"].cast("string"))
    
    import xml.etree.ElementTree as ET
    import json
    
    def returnV(col):
      elem_dict= {}
      tag_list = [
        './TAG/Document/id',
        './TAG/Document/car',
        './TAG/Document/motor',
        './Metadata/Date']
      
      root = ET.fromstring(col)
      
      for tag in tag_list:
        for item in root.findall(tag):
          elem_dict[item.tag] = item.text
      return json.dumps(elem_dict)
    

    我有一些嵌套的 TAG,通过这种方法,我提取了所有需要的值并将它们作为 JSON 返回。我了解到的是,如果传入架构可以更改,结构化流不是解决方案。所以我只取了那些我知道它们不会随着时间而改变的值。

    然后,一旦定义了方法,我就将它注册为 UDF。

    extractValuesFromXML = udf(returnV)
    XML_DF= df.withColumn("body",extractValuesFromXML("body"))
    

    最后我只是使用get_json_object函数来提取JSON的值

    input_parsed_df = XML_DF.select(
      get_json_object("body", "$.id").alias("id").cast('integer'), 
      get_json_object("body", "$.car").alias("car"),
      get_json_object("body", "$.motor").alias("motor"),
      get_json_object("body", "$.Date").alias("Date")
    
    )
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-06-22
      • 1970-01-01
      • 1970-01-01
      • 2021-09-16
      • 1970-01-01
      相关资源
      最近更新 更多