I am converting some hive tabular data into JSON documents using pyspark and writing output to HDFS for downstream consumption. code snippt looks like this:
def convertToJson(x):
data = ....... #transformation_code
output = row(str(json.dumps(data))
return output
df1 = spark.sql("""select * from json_ready_tbl""")
rdd1 = df.rdd.map(lambda x: convertToJson(x)).saveAsTextFile('/hdfs/output/dir')
The output JSON of convertToJSON function is nested but at a high level looks something like this:
{
"id": 1001,
"info": {
"count": 12345,
"code": 999
}
}
So basically, the output files have a JSON string in each row. We then create a hive table on top of this text files by adding some metadata to each row like date and client_name for further downstream consumption as follows:
date client_name json_obj
2020-04-01 zeus { "id": 1001, "info": { "count": 12345, "code": 999 } }
...
N number of rows
Now the challenge is, we have to batch these JSON files upto 1000 into an array of JSONs and write to HDFS so the output in hive table should look like this:
date client_name json_batch
2020-04-01 zeus [{obj1}, {obj2}, {obj3},...{obj1000}]
...
T number of rows
How can I achieve this? Any help is appreciated. Thanks.