Estoy usando AWS para transformar algunos archivos JSON. He agregado los archivos a Glue desde S3. El trabajo que configuré lee los archivos en ok, el trabajo se ejecuta correctamente, hay un archivo agregado al depósito S3 correcto. El problema que tengo es que no puedo nombrar el archivo: se le asigna un nombre aleatorio, tampoco se le asigna la extensión .JSON.
¿Cómo puedo nombrar el archivo y también agregar la extensión a la salida?
Debido a la naturaleza del funcionamiento de Spark, no es posible nombrar el archivo. Sin embargo, es posible cambiar el nombre del archivo inmediatamente después.
URI = sc._gateway.jvm.java.net.URI Path = sc._gateway.jvm.org.apache.hadoop.fs.Path FileSystem = sc._gateway.jvm.org.apache.hadoop.fs.FileSystem fs = FileSystem.get(URI("s3://{bucket_name}"), sc._jsc.hadoopConfiguration()) file_path = "s3://{bucket_name}/processed/source={source_name}/year={partition_year}/week={partition_week}/" df.coalesce(1).write.format("json").mode( "overwrite").option("codec", "gzip").save(file_path) # rename created file created_file_path = fs.globStatus(Path(file_path + "part*.gz"))[0].getPath() fs.rename( created_file_path, Path(file_path + "{desired_name}.jl.gz"))Este siguiente código funcionó para mí:
source_DataFrame = glueContext.create_dynamic_frame.from_catalog(database = databasename, table_name = source_tablename_in_catalog, transformation_ctx = "source_DataFrame") source_DataFrame = source_DataFrame.toDF().coalesce(1) #avoiding coalesce(1) will create many part-000* files according to data from awsglue.dynamicframe import DynamicFrame DyF = DynamicFrame.fromDF(source_DataFrame, glueContext, "DyF") # writing the file as usual in Glue. **I have given some partitions** too. # keep "partitionKeys":[] in case of no partitions output_Parquet = glueContext.write_dynamic_frame.from_options(frame = DyF, connection_type = "s3", format = "parquet", connection_options = {"path": destination_path + "/", "partitionKeys": ["department","team","card","datepartition"]}, transformation_ctx = "output_Parquet") import boto3 client = boto3.client('s3') #getting all the content/file inside the bucket. response = client.list_objects_v2(Bucket=bucket_name) names = response["Contents"] #Find out the file which have part-000* in it's Key particulars = [name['Key'] for name in names if 'part-000' in name['Key']] #Find out the prefix of part-000* because we want to retain the partitions schema location = [particular.split('part-000')[0] for particular in particulars] #Constrain - copy_object has limit of 5GB.datepartition=20190131 for key,particular in enumerate(particulars): client.copy_object(Bucket=bucket_name, CopySource=bucket_name + "/" + particular, Key=location[key]+"newfile") client.delete_object(Bucket=bucket_name, Key=particular) job.commit()La piedra angular es que fallará al copiar el archivo (copy_object) cuando tenga más de 5 GB. Puedes usar esto
s3 = boto3.resource('s3') for key,particular in enumerate(particulars): copy_source = { 'Bucket': bucket_name, 'Key': particular } s3.meta.client.copy(copy_source, bucket_name, location[key]+"newfile")