For a specific use case, I need to read data from a delta lake I have on Azure (stored on adls gen2) and write the data to a JSONL file (one json per line) that is in blob storage. I understand that converting to pandas and writing to a single json file means I am not benefitting pyspark scaling, but for this use case I need to write to a single json that is named based on the date (once written to the json, the json is ~700MB to 1GB, so it is a large file).
What I have been doing is reading the data in a pyspark notebook on Databricks, formatting it in the format I want by creating a few columns that are structs in order to get a nested column, converting it to a pandas dataframe, and writing the jsonl on azure storage using the azure-storage python sdk.
The problem I am having with this method is that I have some structs, and while the data looks good when I display it, when I write it the nested column names are missing. Here is an example of 1 line of the json lines file (there are tens of thousands of lines). For simplicity, I removed a lot of the fields. The SENSOR_INFO nested field also has additional fields.
{"DATE": "2022-01-02", "LOCATION":"XX", ..., "SENSOR_INFO":[1292900, 231, 1.2]}
Expected output
{"DATE": "2022-01-02", "LOCATION":"XX", ..., "SENSOR_INFO": {"ID": 1292900, "NUM_PINGS": 231, "CALC_1": 1.2}}
Each file is written and must be saved with a specific agreed upon name based on the date and can be overwritten if there is a backfill or new data for a date comes in. This process runs frequently and could be running for 2 dates at once. All of the files are in one directory and not nested within a container.
Here is the code I am using:
from pyspark.sql import functions as F
from azure.storage.blob import BlobServiceClient # use this to write, have azure-storage installed on cluster
spark.conf.set("spark.sql.execution.arrow.enabled", "false") # has to be set to false for structs to convert?
# Load the data from delta table/path
df = spark \
.read \
.format('delta') \
.load(load_path)
df_export = df \
.withColumn("SENSOR_INFO", F.struct(F.col("SENSOR_ID").alias("ID"), F.col("SENSOR_NUM_PINGS").alias("NUM_PINGS"), F.col("SENSOR_CALC_1").alias("CALC_1"))) \
.drop("SENSOR_ID", "SENSOR_PINGS", "SENSOR_CALC_1")
pandas_df_export = df_export.select("*").toPandas()
json_data = pandas_df_export.to_json(orient='records', lines=True, double_precision=2)
# blob_name includes the date, it is an agreed upon file naming format .jsonl
blob_service_client = BlobServiceClient.from_connection_string(connect_str)
blob_client = blob_service_client.get_blob_client(container=container_name, blob=blob_name)
blob_client.upload_blob(json_data, overwrite=True)
Up until a few months ago, this process was working and producing the expected output. But then something must have changed, and all of the sudden it started failing and I had to switch spark.sql.execution.arrow.enabled to false for it to even work. I am not sure if that is what broke it, or if it was something else.
Unfortunately, I cannot change the file format to remove the nested field as it will be a breaking change downstream.
I do not need to use structs, I just want nested fields/names. I have read that one way to solve this is to do create_map instead of struct, but that only works if all of the fields in the map are the same type. I have also read about compacting files to a single file in delta lake, but because this process runs frequently I am not sure if that makes sense... It also means the order of the data in the file will likely be somewhat random based on the different files being combined (not the most important thing, but not ideal [there is a line where I sort by a few columns that is not included]).
I do not need to export the data using a Databricks notebook given I am not using pyspark to write, if there is a better way to do it I am open to doing it elsewhere. Due to the fact that I am reading from delta lake I didn't know of a better way as exporting data to a single file is something I need to do frequently for external purposes. I am open to using an Azure function or some better mechanism if it exists. This process is already very slow so ideally if there is something more efficient, that would be even better.