Add File Modified Column to JSON Parsing Transform

Viewed 58

I have a transform that takes in .json.gz files as input, there’s a large number of different json schemas that I’m writing out to different outputs, so I’m hoping I can infer the schema. So far, I’ve had success using spark.read.json(paths), however I’ve come to realize I need to add a column that specifies the FileStatus.modified timestamp as a column in the output dataset for downstream transform purposes.

It appears this is possible using rdd.flatMap(process_file) similar to transforms.verbs.files.json_to_df (this only supports .json not .json.gz). I could define a pass a custom process_file function to rdd.flatMap that unzips the .gz, parses the json, and attempts to infer schema - however I lose the robustness of using spark.read.json(paths).

Any better options make sense here? I’m surprised this isn’t supported in spark.read.json(paths) but I could be missing something.

1 Answers

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.
Related