From 24e02b81d009c662fb8bf03b91fb63d6a660bd8a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20K=C4=99drzy=C5=84ski?= Date: Thu, 19 May 2022 13:26:35 +0200 Subject: [PATCH 1/6] fix: my little change (#1) --- composer/cicd_sample/dags/example_dag.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/composer/cicd_sample/dags/example_dag.py b/composer/cicd_sample/dags/example_dag.py index f147ef7580d..07bae8ee5f5 100644 --- a/composer/cicd_sample/dags/example_dag.py +++ b/composer/cicd_sample/dags/example_dag.py @@ -16,6 +16,8 @@ from airflow import models from airflow.operators import bash +# foobar + # If you are running Airflow in more than one time zone # see https://airflow.apache.org/docs/apache-airflow/stable/timezone.html # for best practices From 1fe9ad4ec81e6b0965324efec4449eb0b1970efb Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20K=C4=99drzy=C5=84ski?= Date: Thu, 19 May 2022 13:44:28 +0200 Subject: [PATCH 2/6] fix: another one (#2) --- composer/cicd_sample/dags/example_dag.py | 1 + 1 file changed, 1 insertion(+) diff --git a/composer/cicd_sample/dags/example_dag.py b/composer/cicd_sample/dags/example_dag.py index 07bae8ee5f5..223f9b721f1 100644 --- a/composer/cicd_sample/dags/example_dag.py +++ b/composer/cicd_sample/dags/example_dag.py @@ -17,6 +17,7 @@ from airflow.operators import bash # foobar +# bar # If you are running Airflow in more than one time zone # see https://airflow.apache.org/docs/apache-airflow/stable/timezone.html From be818c3164e1c6a92700763735499b4e111d49a4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20K=C4=99drzy=C5=84ski?= Date: Thu, 19 May 2022 13:49:19 +0200 Subject: [PATCH 3/6] fix: one more --- composer/cicd_sample/dags/example_dag.py | 1 - 1 file changed, 1 deletion(-) diff --git a/composer/cicd_sample/dags/example_dag.py b/composer/cicd_sample/dags/example_dag.py index 223f9b721f1..07bae8ee5f5 100644 --- a/composer/cicd_sample/dags/example_dag.py +++ b/composer/cicd_sample/dags/example_dag.py @@ -17,7 +17,6 @@ from airflow.operators import bash # foobar -# bar # If you are running Airflow in more than one time zone # see https://airflow.apache.org/docs/apache-airflow/stable/timezone.html From fac212e00d148195be3bf286886fa74bda5672e4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20K=C4=99drzy=C5=84ski?= Date: Thu, 19 May 2022 13:50:26 +0200 Subject: [PATCH 4/6] fix: asdf --- composer/cicd_sample/dags/example_dag.py | 1 + 1 file changed, 1 insertion(+) diff --git a/composer/cicd_sample/dags/example_dag.py b/composer/cicd_sample/dags/example_dag.py index 07bae8ee5f5..f7d2ddb07f5 100644 --- a/composer/cicd_sample/dags/example_dag.py +++ b/composer/cicd_sample/dags/example_dag.py @@ -17,6 +17,7 @@ from airflow.operators import bash # foobar +# foo # If you are running Airflow in more than one time zone # see https://airflow.apache.org/docs/apache-airflow/stable/timezone.html From 7eaed522249c0c39557f17f6eb4b349124c010c8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20K=C4=99drzy=C5=84ski?= Date: Thu, 19 May 2022 13:56:09 +0200 Subject: [PATCH 5/6] Update add-dags-to-composer.cloudbuild.yaml --- composer/cicd_sample/add-dags-to-composer.cloudbuild.yaml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/composer/cicd_sample/add-dags-to-composer.cloudbuild.yaml b/composer/cicd_sample/add-dags-to-composer.cloudbuild.yaml index 3a1a6a6b68c..3820ed0f9ab 100644 --- a/composer/cicd_sample/add-dags-to-composer.cloudbuild.yaml +++ b/composer/cicd_sample/add-dags-to-composer.cloudbuild.yaml @@ -17,11 +17,11 @@ steps: # install dependencies - name: python entrypoint: pip - args: ["install", "-r", "utils/requirements.txt", "--user"] + args: ["install", "-r", "composer/cicd_sample/utils/requirements.txt", "--user"] # run - name: python entrypoint: python - args: ["utils/add_dags_to_composer.py", "--dags_directory=${_DAGS_DIRECTORY}", "--dags_bucket=${_DAGS_BUCKET}"] + args: ["composer/cicd_sample/utils/add_dags_to_composer.py", "--dags_directory=${_DAGS_DIRECTORY}", "--dags_bucket=${_DAGS_BUCKET}"] # [END composer_cicd_dagsync_yaml] From d86fde9b86cb3c3ab20c3e4b3df328da1f29f87a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20K=C4=99drzy=C5=84ski?= Date: Thu, 19 May 2022 14:36:49 +0200 Subject: [PATCH 6/6] feat: added data analytics dag --- .../cicd_sample/dags/data_analytics_dag.py | 125 ++++++++++++++++++ 1 file changed, 125 insertions(+) create mode 100644 composer/cicd_sample/dags/data_analytics_dag.py diff --git a/composer/cicd_sample/dags/data_analytics_dag.py b/composer/cicd_sample/dags/data_analytics_dag.py new file mode 100644 index 00000000000..61004603a96 --- /dev/null +++ b/composer/cicd_sample/dags/data_analytics_dag.py @@ -0,0 +1,125 @@ +# Copyright 2022 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + + +import datetime + +from airflow import models +from airflow.providers.google.cloud.operators import dataproc +from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator +from airflow.providers.google.cloud.transfers.gcs_to_bigquery import ( + GCSToBigQueryOperator, +) +from airflow.utils.task_group import TaskGroup + +PROJECT_NAME = "{{var.value.gcp_project}}" + +# BigQuery configs +BQ_DESTINATION_DATASET_NAME = "holiday_weather_michal" +BQ_DESTINATION_TABLE_NAME = "holidays_weather_michal_joined" +BQ_NORMALIZED_TABLE_NAME = "holidays_weather_michal_normalized" + + +PYSPARK_JAR = "gs://spark-lib/bigquery/spark-bigquery-latest_2.12.jar" +PROCESSING_PYTHON_FILE = "gs://{{var.value.gcs_bucket}}/data_analytics_process.py" + +BATCH_ID = "test-batch-id-{{ ts_nodash | lower}}" # Dataproc serverless only allows lowercase characters +BATCH_CONFIG = { + "pyspark_batch": { + "jar_file_uris": [PYSPARK_JAR], + "main_python_file_uri": PROCESSING_PYTHON_FILE, + "args": [ + PROJECT_NAME, + f"{BQ_DESTINATION_DATASET_NAME}.{BQ_DESTINATION_TABLE_NAME}", + f"{BQ_DESTINATION_DATASET_NAME}.{BQ_NORMALIZED_TABLE_NAME}", + ], + }, +} + +yesterday = datetime.datetime.combine( + datetime.datetime.today() - datetime.timedelta(1), datetime.datetime.min.time() +) + +default_dag_args = { + # Setting start date as yesterday starts the DAG immediately when it is + # detected in the Cloud Storage bucket. + "start_date": yesterday, + # To email on failure or retry set 'email' arg to your email and enable + # emailing here. + "email_on_failure": False, + "email_on_retry": False, + # If a task fails, retry it once after waiting at least 5 minutes + "retries": 1, + "retry_delay": datetime.timedelta(minutes=5), +} + +with models.DAG( + "summit_dag", + # Continue to run DAG once per day + schedule_interval=datetime.timedelta(days=1), + default_args=default_dag_args, +) as dag: + + create_batch = dataproc.DataprocCreateBatchOperator( + task_id="create_batch", + project_id=PROJECT_NAME, + region="{{ var.value.gce_region }}", + batch=BATCH_CONFIG, + batch_id=BATCH_ID, + ) + load_external_dataset = GCSToBigQueryOperator( + task_id="run_bq_external_ingestion", + bucket="{{var.value.gcs_bucket}}", + source_objects=["holidays.csv"], + destination_project_dataset_table=f"{BQ_DESTINATION_DATASET_NAME}.holidays", + source_format="CSV", + schema_fields=[ + {"name": "Date", "type": "DATE"}, + {"name": "Holiday", "type": "STRING"}, + ], + skip_leading_rows=1, + ) + + with TaskGroup("join_bq_datasets") as bq_join_group: + + for year in range(1997, 2022): + # BigQuery configs + BQ_DATASET_NAME = f"bigquery-public-data.ghcn_d.ghcnd_{str(year)}" + BQ_DESTINATION_TABLE_NAME = "holidays_weather_joined" + # Specifically query a Chicago weather station + WEATHER_HOLIDAYS_JOIN_QUERY = f""" + SELECT Holidays.Date, Holiday, id, element, value + FROM `{PROJECT_NAME}.holiday_weather.holidays` AS Holidays + JOIN (SELECT id, date, element, value FROM {BQ_DATASET_NAME} AS Table WHERE Table.element="TMAX" AND Table.id="USW00094846") AS Weather + ON Holidays.Date = Weather.Date; + """ + + bq_join_holidays_weather_data = BigQueryInsertJobOperator( + task_id=f"bq_join_holidays_weather_data_{str(year)}", + configuration={ + "query": { + "query": WEATHER_HOLIDAYS_JOIN_QUERY, + "useLegacySql": False, + "destinationTable": { + "projectId": PROJECT_NAME, + "datasetId": BQ_DESTINATION_DATASET_NAME, + "tableId": BQ_DESTINATION_TABLE_NAME, + }, + "writeDisposition": "WRITE_APPEND", + } + }, + location="US", + ) + + load_external_dataset >> bq_join_group >> create_batch