I assume you have a dataset with one or multiple JSONs files gzipped. Each JSON file has 1 json or multiple json(s), one per line.
Both solutions could be adapted to multilines JSONs.
Your current implementation
You can pass a list of paths to spark.read.json(paths), so you likely generate a list of paths.
Your transform probably looks like :
@transform(
out=Output("your_output"),
source=Input("your_input")
)
def parse_jsons(ctx, source, out):
# Generate a list of paths
filesystem = source.filesystem()
hadoop_path = filesystem.hadoop_path
paths = [f"{hadoop_path}/{f.path}" for f in filesystem.ls()]
# Parse those files with spark built-ins.
df = ctx.spark_session.read.json(paths)
out.write_dataframe(df)
Example solution
You can set a schema on the input dataset to let Foundry parse each json in its own "cell" and then apply the spark built-ins on each cell. The main difference is that it operates at a dataframe level, which should bring a tiny advantage in your context : being able to add columns.
{
"fieldSchemaList": [
{
"type": "STRING",
"name": "row",
"nullable": null,
"userDefinedTypeClass": null,
"customMetadata": {},
"arraySubtype": null,
"precision": null,
"scale": null,
"mapKeyType": null,
"mapValueType": null,
"subSchemas": null
}
],
"primaryKey": null,
"dataFrameReaderClass": "com.palantir.foundry.spark.input.DataSourceDataFrameReader",
"customMetadata": {
"format": "text",
"options": {}
}
}
Your transform will then apply the parsing logic on each cell :
@transform(
out=Output("your_output"),
source_df=Input("your_input"),
)
def compute(ctx, source_df, out):
# Get the file system
filesystem = source_df.filesystem().files()
# Convert to df and create a path column
source_df = source_df.dataframe().withColumn("full_path", F.input_file_name())
# Generate Schema
schema = ctx.spark_session.read.json(source_df.rdd.map(lambda x: x["row"])).schema
# Parse JSONs with the discovered schema
results = source_df.withColumn("json_parsed", F.from_json("row", schema))
# Convert the paths into filenames and join with filesystem (to get modified timestamp)
results = results.withColumn("path", F.reverse(F.split("full_path", "/"))[0])
results = results.join(filesystem, "path", "left")
# Extract all columns
results = results.select(*results.columns, "json_parsed.*")
out.write_dataframe(results)
Areas of improvement :
- File name obtained via input_file_name are URL escaped (
my json.json shows as my%20json.json). If you files have special characters, you will need to decode them. See Decoding a String URL column in pyspark?
- The above snippet won't generate a correct output in Preview mode as temporary files are used in this case (hence the file names of the
filesystem() and input_file_name() won't match).
Notes :
- those methods should supports a mix of gzipped and non-gzipped JSONs in the same dataset, given each file is parsed independently and spark handle both.
- JSON which don't have a particular key will generate a "null" value in the output column of the missing key.
- You also get the file size, via this method, as it is available in the filesystem() dataframe and part of the join.