import sys from awsglue.utils import getResolvedOptions import boto3 from utils.util import ( init_spark_context, get_environment, get_data_from_postgresql_and_save_to_s3, ) args = getResolvedOptions(sys.argv, ["JOB_NAME", "partition_date"]) # daily p_date = args["partition_date"] table_name = "bots_usageplan" print(f"p_date: {p_date}") sc, glueContext, spark, job = init_spark_context() environment = get_environment() if environment == "PROD": s3_bucket = "analytics-database-data" elif environment == "STAGING": s3_bucket = "analytics-database-data-staging" else: raise Exception("Invalid environment.") sample_query = f"SELECT * FROM {table_name}" transform_sql = f""" select *, '{p_date}' as pdate from result """ # Clean up existing target prefix before writing s3_prefix = f"{table_name}/pdate={p_date}" s3 = boto3.resource("s3") bucket_obj = s3.Bucket(s3_bucket) to_delete = list(bucket_obj.objects.filter(Prefix=s3_prefix)) if to_delete: print(f"Cleaning up {len(to_delete)} existing objects under s3://{s3_bucket}/{s3_prefix}/ ...") chunk = [] for obj in to_delete: chunk.append({"Key": obj.key}) if len(chunk) == 1000: s3.meta.client.delete_objects(Bucket=s3_bucket, Delete={"Objects": chunk}) chunk = [] if chunk: s3.meta.client.delete_objects(Bucket=s3_bucket, Delete={"Objects": chunk}) else: print(f"No existing objects found under s3://{s3_bucket}/{s3_prefix}/") row_count = get_data_from_postgresql_and_save_to_s3( glueContext=glueContext, spark=spark, transform_sql=transform_sql, db_connection_options={ "useConnectionProperties": "true", "dbtable": table_name, "connectionName": "analystic-database-connection", "sampleQuery": sample_query, }, s3_connection_options={ "path": f"s3://{s3_bucket}/{table_name}/", "partitionKeys": ["pdate"], }, ) print(f"Wrote {row_count} rows to s3://{s3_bucket}/{table_name}/pdate={p_date}/") job.commit()