- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
10-14-2021 02:23 AM
Many Thanks for your response HubertDudek. As mentioned in response, please find the following code which I am using:
*************************************
import os
import pyspark
from pyspark.sql import SparkSession
from pyspark.sql.types import StructType
from pyspark.sql.types import ArrayType
from pyspark.sql.functions import col
from pyspark.sql.functions import explode_outer
from array import array
from azure.storage.blob import BlockBlobService
from datetime import date, timedelta
block_blob_service = BlockBlobService(account_name="********", account_key="*************")
containers = block_blob_service.list_containers()
for c in containers:
top_level_container_name = c.name
generator = block_blob_service.list_blobs(top_level_container_name)
#print(c.name)
if "self-verification" in c.name:
for blob in generator:
if "/PageViews/" in blob.name:
if (date.today() - timedelta(1)).isoformat() in blob.name:
#print(c.name)
df2 = spark.read.option("multiline","true").option("inferSchema","true").option("header","True") .option("recursiveFileLookup","true").json("/mnt/"+c.name+"/"+blob.name)
#print(df2)
def Flatten(df2):
complex_fields = dict([(field.name, field.dataType)
for field in df2.schema.fields
if type(field.dataType) == ArrayType or type(field.dataType) == StructType])
while len(complex_fields) != 0:
col_name = list(complex_fields.keys())[0]
if (type(complex_fields[col_name]) == StructType):
expanded = [col(col_name + '.' + k).alias(col_name + '_' + k) for k in [ n.name for n in complex_fields[col_name]]]
# print(col_name)
# display(df2)
# print(expanded)
df2 = df2.select("*", *expanded).drop(col_name)
#print(df2)
elif (type(complex_fields[col_name]) == ArrayType):
df2 = df2.withColumn(col_name, explode_outer(col_name))
complex_fields = dict([(field.name, field.dataType)
for field in df2.schema.fields
if type(field.dataType) == ArrayType or type(field.dataType) == StructType])
#return df
#print("good")
#print("morning")
return df2
Flatten_df2 = Flatten(df2)
Flatten_df2.write.mode("append").json("/usr/hive/warehouse/stg_pageviews")
************************************
Please help me on this. Many Thanks