Process from Pyspark load from GCS, transform and save to a different GCS not working

Viewed 724

I'm trying to utilize Dataproc (using Pyspark) to load a large dataset form GCS, transform it with geospatial enrichment, and then save back in a PartitionBy format.

So first the code used in airflow to spin up the Dataproc Cluster

from datetime import timedelta, datetime
from airflow import models
from airflow.operators.bash_operator import BashOperator
from airflow.contrib.operators.gcs_to_bq import GoogleCloudStorageToBigQueryOperator

from airflow.contrib.operators import dataproc_operator
from airflow.utils import trigger_rule

# Add 5 day buffer due to UM Lag (i.e. To look back 2 weeks/14 days its 19 instead of 14
yesterday = datetime(2021, 10, 30)
PROJECT = {our_project}

# Spark references
SPARK_CODE = ('{bucket}/spark_script.py') #####CHANGE
dataproc_job_name = 'CAN-Data-Tests' #### CHANGE

default_args = {
    'owner': 'TESTING_FIRST',
    'depends_on_past': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
    'start_date': yesterday,
}

spark_config = {'spark:spark.driver.maxResultSize': '2G', 'spark:spark.kryoserializer.buffer.max': '2040m',
                "spark:spark.task.maxDirectResultSize": "2G", "spark:spark.driver.memoryOverhead'": '5G',
                'spark:spark.serializer': 'org.apache.spark.serializer.KryoSerializer',
                "spark:spark.kryo.registrator": "org.apache.sedona.core.serde.SedonaKryoRegistrator",
                'distcp:mapreduce.map.java.opts': '-Xmx2304m', 'distcp:mapreduce.map.memory.mb': '3024',
                'distcp:mapreduce.reduce.java.opts': '-Xmx2304m', 'distcp:mapreduce.reduce.memory.mb': '3024',
                'spark:spark.executor.cores': '3', 'spark:spark.shuffle.memoryFraction': '0.5',
                'spark:spark.executor.memory': '30G', 'spark:spark.driver.memory': '30G',
                'spark:spark.network.timeout': '6000s'
                }

with models.DAG('Initial_Canada_Testing', description='Geospark', catchup=False, default_args=default_args, schedule_interval=None) as dag1:

    # Create small dataproc cluster
    create_dataproc = dataproc_operator.DataprocClusterCreateOperator(
        task_id='create_dataproc',
        cluster_name='dataproc-cluster-{{ ds_nodash }}',
        project_id=PROJECT,
        region='us-central1',
        num_workers=20,
        zone='us-central1-a',
        image_version='2.0-debian10',
        properties=spark_config,
        master_disk_size=1000,
        worker_disk_size=1000,
        use_if_exists=True,
        init_actions_uris=["gs://{bucket}/data/data_pandas.sh"],
        autoscaling_policy=f'projects/{PROJECT}/locations/us-central1/autoscalingPolicies/autoscale-simple',
        master_machine_type='n2-standard-8',
        worker_machine_type='n2-highmem-16')

    # Run the PySpark job
    run_spark = dataproc_operator.DataProcPySparkOperator(
        task_id='Run_Load_Job',
        main=SPARK_CODE,
        cluster_name='dataproc-cluster-{{ ds_nodash }}',
        job_name=dataproc_job_name,
        region='us-central1',
        dataproc_jars=[
            "gs://bucket/data/gcs-connector-hadoop3-latest.jar",
            "gs://spark-lib/bigquery/spark-bigquery-latest_2.12.jar",
            'gs://bucket/data/sedona-python-adapter-3.0_2.12-1.0.1-incubating.jar',
            'gs://bucket/data/sedona-sql-3.0_2.12-1.0.1-incubating.jar',
            'gs://bucket/data/geotools-wrapper-geotools-24.1.jar',
            "gs://bucket/data/spark-bigquery-with-dependencies_2.12-0.21.1.jar",
            "gs://bucket/data/gcs-connector-hadoop2-2.1.3-shaded.jar",
            'gs://bucket/data/sedona-viz-3.0_2.12-1.0.1-incubating.jar',

        ])

As for the actual Spark Code, this is what I have:

import threading
import queue
import re
import geopandas as gpd
import time
import os
import pandas as pd
import pytz
import shapely
from datetime import datetime
from datetime import timedelta, date
from google.cloud import storage

from pyspark.sql import SparkSession
from pyspark import SparkContext
from pyspark import StorageLevel
from pyspark.sql import SparkSession
from pyspark.sql.functions import when, lit, col, year, month, dayofmonth
from pyspark.sql.functions import from_unixtime
from pyspark.sql.types import StructField, StructType, StringType, IntegerType, DoubleType, LongType
from pyspark.sql.functions import current_date
from pyspark.sql.functions import col, unix_timestamp, to_date, from_unixtime

from sedona.utils.adapter import Adapter
from sedona.register import SedonaRegistrator
from sedona.utils import KryoSerializer, SedonaKryoRegistrator


spark = SparkSession.builder \
    .config("spark.jars.packages", "com.google.cloud.bigdataoss:gcs-connector:hadoop3-2.2.0") \
    .config("spark.jars.packages", "com.google.cloud.spark:spark-bigquery-with-dependencies_2.12:0.21.1") \
    .config('spark.jars.packages', 'org.apache.sedona:sedona-python-adapter-3.0_2.12:1.0.0-incubating') \
    .config('spark.jars.packages', 'org.datasyslab:geotools-wrapper:geotools-24.0') \
    .config("spark.serializer", KryoSerializer.getName) \
    .config("spark.kryo.registrator", SedonaKryoRegistrator.getName) \
    .getOrCreate()

spark._jsc.hadoopConfiguration().set('fs.gs.impl', 'com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem')
spark._jsc.hadoopConfiguration().set('fs.gs.auth.service.account.enable', 'true')
spark.conf.set("google.cloud.auth.service.account.enable", "true")

SedonaRegistrator.registerAll(spark)
conf = spark.sparkContext._jsc.hadoopConfiguration()

conf.set("spark.executor.memory", "35g")
conf.set("spark.driver.memory", "35")
conf.set("spark.spark.executor.cores", '3')
conf.set('spark.sql.session.timeZone', 'US/Central')
# conf.set("spark.executor.instances", '5')
# conf.set("spark.cores.max", '15')

# Create a schema for the dataframe
schema = StructType([
    StructField("Hashed_Device_ID", StringType(), True), StructField("Unix_Timestamp_of_Visit", LongType(), True),
    StructField("Lat_of_Visit", DoubleType(), True), StructField("Lon_of_Visit", DoubleType(), True),
    StructField("Accuracy", DoubleType(), True), StructField("Datasource_ID", StringType(), True),
    StructField("ip", StringType(), True), StructField("gps_source", StringType(), True),
    StructField("country", StringType(), True), StructField("carrier", StringType(), True),
    StructField("app_name", StringType(), True), StructField("app_id", StringType(), True),
    StructField("ua", StringType(), True), StructField("os", StringType(), True),
    StructField("os_version", StringType(), True), StructField("categories", StringType(), True),
    StructField("keywords", StringType(), True), StructField("publisher_id", StringType(), True),
    StructField("publisher_name", StringType(), True), StructField("is_hashed", StringType(), True),
    StructField("make", StringType(), True), StructField("model", StringType(), True),
    StructField("id_type", StringType(), True), StructField("gender", StringType(), True),
    StructField("birth_year", IntegerType(), True), StructField("gps_country", StringType(), True),
    StructField("heading", DoubleType(), True), StructField("speed", DoubleType(), True),
    StructField("altitude", DoubleType(), True), StructField("vertical_accuracy", DoubleType(), True),
    StructField("gps_subsource", StringType(), True), StructField("wifi_bssid", StringType(), True),
    StructField("wifi_ssid", StringType(), True), StructField("imei", StringType(), True),
    StructField("imsi", StringType(), True)
])


test_file_path = 'gs://data_bucket/region=canada/year=2019/month=01/day=01/part-00000-5a5f143a-e5e8-4f5d-9a62-afbf0020cdf3.c000.csv.gz'

raw_data = spark.read.csv(test_file_path, schema=schema, sep='\t', ignoreLeadingWhiteSpace=True, ignoreTrailingWhiteSpace=True)
raw_data.createOrReplaceTempView("raw_table")

bq_prov = spark.read \
  .format("bigquery") \
  .option("table", "UberMedia_Canada_All.Canadian_Provinces") \
  .load()

bq_prov.cache()

bq_prov.createOrReplaceTempView("bq_provinces")

final_df = spark.sql(
    """
    WITH raw_data_table AS(
        SELECT *, ST_Point(Lon_of_Visit, Lat_of_Visit) AS pt_geom
        FROM raw_table
    ),
    
    bq_provs AS(
        SELECT name, ST_GeomFromWKT(geometry) as geom
        FROM bq_provinces       
    )
    
    SELECT a.*, b.name as province
    FROM raw_data_table as a
    LEFT OUTER JOIN bq_provs as b
    ON ST_Contains(b.geom, a.pt_geom)
    """
)


final_df = final_df.drop("pt_geom")
final_df = final_df.withColumn("Unix_Timestamp_of_Visit", from_unixtime(col("Unix_Timestamp_of_Visit")))
final_df = final_df.withColumn("partition_prov", col("province"))
final_df = final_df.withColumn("year", year(to_date(col("Unix_Timestamp_of_Visit"))))
final_df = final_df.withColumn("month", month(to_date(col("Unix_Timestamp_of_Visit"))))
final_df = final_df.withColumn("day", dayofmonth(to_date(col("Unix_Timestamp_of_Visit"))))



final_df.write \
    .partitionBy('prov','year','month','day') \
    .parquet('gs://save_bucket/canada/', mode="overwrite", compression="snappy")

So let me preface, by I'm still SUPER new at Clusters and Dataproc. This initial testing was done on a small sample (one day of data). I want end result to parition by province,year, month, and then day. Again I've ran this solely on one day, but I have about 4 years worth of data. My initial expectations where it would run relatively quickly, but I was sadly mistaken. I kicked of the job and it would be very active initially and then it would hang up with logs simply repeating:

[
  {
    "insertId": "pb2dx292npp3us6qq",
    "jsonPayload": {
      "class": "org.apache.hadoop.mapreduce.v2.hs.JobHistory",
      "message": "Starting scan to move intermediate done files",
      "filename": "hadoop-mapred-historyserver-dataproc-cluster-demo-20210927-m.log"
    },
    "resource": {
      "type": "cloud_dataproc_cluster",
      "labels": {
        "region": "us-central1",
        "project_id": "project",
        "cluster_name": "dataproc-cluster-demo-20210927",
        "cluster_uuid": "2291cf56-9555-4ef5-93b4-ea6f5793fcc1"
      }
    },
    "timestamp": "2021-11-03T08:08:19.440Z",
    "severity": "INFO",
    "labels": {
      "compute.googleapis.com/zone": "us-central1-a",
      "compute.googleapis.com/resource_id": "2070420762412198322",
      "compute.googleapis.com/resource_name": "dataproc-cluster-demo-20210927-m"
    },
    "logName": "projects/project/logs/hadoop-mapred-historyserver",
    "receiveTimestamp": "2021-11-03T08:08:23.636556798Z"
  },

Any help would appreciated and can add other info needed.

Thanks again!

1 Answers

Please change you SparkSession building so that all the packages are in the same call:

spark = SparkSession.builder \
    .config("spark.jars.packages", 'com.google.cloud.bigdataoss:gcs-connector:hadoop3-2.2.0,com.google.cloud.spark:spark-bigquery-with-dependencies_2.12:0.21.1,org.apache.sedona:sedona-python-adapter-3.0_2.12:1.0.0-incubating,org.datasyslab:geotools-wrapper:geotools-24.0') \
    .config("spark.serializer", KryoSerializer.getName) \
    .config("spark.kryo.registrator", SedonaKryoRegistrator.getName) \
    .getOrCreate()

See more at the spark configuration documentation

Related