Spark Structured Streaming Dynamic Parsing of XML

Viewed 172

Currently im consuming xml data from a kafka source and process them with Spark Structured Streaming.

To get the needed information out of the xml i am using xpath. As i want to make the pipeline more dynamic i tried to implement a dictionary which hold the column name to be extracted and the expression itself. In a future version the dictionary could get filled by some configuration files without touching the python job.

Unfortunatly it seems not to work as desired (i am a python noob, maybe that's why...).

The xml could be described as followed:

<root>
   <header eventId ="1234" .../>
   <.../>
</root>

My python code looks like this:

df = spark.readStream.format("kafka")...load()
df = df.selectExpr("CAST(timestamp) AS String)", "CAST(value AS String"))
xml_data = df
             .selectExpr("xpath(value, './root/header/@eventId')event_id", ...)
             .selectExpr("explode(arrays_zip(event_id,...)) value"
             .select('value.*')

My next step was defining the dict:

mapping_dict = {
                'event_id' : './root/header/@eventId',
                ...
               }

I tried to rebuild the expression like this:

event_id = "\"xpath(value,'" + mapping_dict.get('event_id) + "') event_id\""

xml_data = df.selectExpr(event_id,...)
             .selectExpr("explode(arrays_zip(event_id,...)) value"
             .select('value.*')

Now i tried to use the dict value in the selectExpr but it fails with an error

org.apache.spark.sql.AnalysisException: cannot resolve '`event_id`' given input columns: [xpath(value, './root/header/@eventId') event_id ...

So this is my first problem, the second one would be, that i want to iterate over this dict and try to extract each entry from the xml. I don't know if i can do that with structured streaming that easily or if i would have to use an udf. And if so, how could an udf for this purpose look like?

Cheers

0 Answers
Related