您可以使用以下代码在单独的行中创建,并将数据写入的单独文件
消息记录
和
消息
。
record = {}
record["field1"] = json_data["field1"]
record["field2"] = json_data["field2"]
message_records_df =spark.createDataFrame([record])
messages_df = spark.createDataFrame([record])
使用创建两个数据帧
field1
和
field2
。由于这两者都与相同
消息记录
和
消息
。
from pyspark.sql.types import LongType
from pyspark.sql import Row
def zipindexdf(df):
schema_new = df.schema.add("index", LongType(), False)
return df.rdd.zipWithIndex().map(lambda l: list(l[0]) + [l[1]]).toDF(schema_new)
message_records_df_index = zipindexdf(message_records_df)
message_records_df_index.show()
messages_df_index = zipindexdf(messages_df)
messages_df_index.show()
我在这里添加
index
列使用
zipWithIndex
所以它可以在索引上联接。
接下来,创建数据帧并通过循环遍历中的每个项合并到最终数据帧
消息记录
如下所示。
for i in json_data['message_records']:
df = zipindexdf(spark.createDataFrame([i]))
message_records_df_index = message_records_df_index.join(df, "index", "inner")
message_records_df_index.show()
和我一样
消息
如下所示
for i in json_data['messages']:
df = zipindexdf(spark.createDataFrame([i]))
messages_df_index = messages_df_index.join(df, "index", "inner")
messages_df_index.show()
最后,将这些数据写入csv文件。
message_records_df_index.write.option("header","true").csv('/message_records_df/')
messages_df_index.write.option("header","true").csv('/messages_df/')
spark.read.option("header","true").csv('/message_records_df/').show()
spark.read.option("header","true").csv('/messages_df/').show()