diff --git a/UPDATING.md b/UPDATING.md index dc47dfe5b103a..031f2ab25d129 100644 --- a/UPDATING.md +++ b/UPDATING.md @@ -41,6 +41,11 @@ assists users migrating to a new version. ## Airflow Master +### Idempotency in BigQuery operators +Idempotency was added to `BigQueryCreateEmptyTableOperator` and `BigQueryCreateEmptyDatasetOperator`. +But to achieve that try / except clause was removed from `create_empty_dataset` and `create_empty_table` +methods of `BigQueryHook`. + ### Migration of AWS components All AWS components (hooks, operators, sensors, example DAGs) will be grouped together as decided in diff --git a/airflow/gcp/hooks/bigquery.py b/airflow/gcp/hooks/bigquery.py index cae5d42a53096..07a56feb5e54b 100644 --- a/airflow/gcp/hooks/bigquery.py +++ b/airflow/gcp/hooks/bigquery.py @@ -329,22 +329,10 @@ def create_empty_table(self, num_retries = num_retries if num_retries else self.num_retries - self.log.info('Creating Table %s:%s.%s', - project_id, dataset_id, table_id) - - try: - self.service.tables().insert( - projectId=project_id, - datasetId=dataset_id, - body=table_resource).execute(num_retries=num_retries) - - self.log.info('Table created successfully: %s:%s.%s', - project_id, dataset_id, table_id) - - except HttpError as err: - raise AirflowException( - 'BigQuery job failed. Error was: {}'.format(err.content) - ) + self.service.tables().insert( + projectId=project_id, + datasetId=dataset_id, + body=table_resource).execute(num_retries=num_retries) def create_external_table(self, # pylint: disable=too-many-locals,too-many-arguments external_project_dataset_table: str, @@ -540,20 +528,14 @@ def create_external_table(self, # pylint: disable=too-many-locals,too-many-argu if encryption_configuration: table_resource["encryptionConfiguration"] = encryption_configuration - try: - self.service.tables().insert( - projectId=project_id, - datasetId=dataset_id, - body=table_resource - ).execute(num_retries=self.num_retries) - - self.log.info('External table created successfully: %s', - external_project_dataset_table) + self.service.tables().insert( + projectId=project_id, + datasetId=dataset_id, + body=table_resource + ).execute(num_retries=self.num_retries) - except HttpError as err: - raise Exception( - 'BigQuery job failed. Error was: {}'.format(err.content) - ) + self.log.info('External table created successfully: %s', + external_project_dataset_table) def patch_table(self, # pylint: disable=too-many-arguments dataset_id: str, @@ -1748,20 +1730,9 @@ def create_empty_dataset(self, dataset_id = dataset_reference.get("datasetReference").get("datasetId") # type: ignore dataset_project_id = dataset_reference.get("datasetReference").get("projectId") # type: ignore - self.log.info('Creating Dataset: %s in project: %s ', dataset_id, - dataset_project_id) - - try: - self.service.datasets().insert( - projectId=dataset_project_id, - body=dataset_reference).execute(num_retries=self.num_retries) - self.log.info('Dataset created successfully: In project %s ' - 'Dataset %s', dataset_project_id, dataset_id) - - except HttpError as err: - raise AirflowException( - 'BigQuery job failed. Error was: {}'.format(err.content) - ) + self.service.datasets().insert( + projectId=dataset_project_id, + body=dataset_reference).execute(num_retries=self.num_retries) def delete_dataset(self, project_id: str, dataset_id: str, delete_contents: bool = False) -> None: """ diff --git a/airflow/gcp/operators/bigquery.py b/airflow/gcp/operators/bigquery.py index a3270608582af..6ec11288d9f6d 100644 --- a/airflow/gcp/operators/bigquery.py +++ b/airflow/gcp/operators/bigquery.py @@ -26,6 +26,8 @@ import warnings from typing import Any, Dict, Iterable, List, Optional, SupportsAbs, Union +from googleapiclient.errors import HttpError + from airflow.exceptions import AirflowException from airflow.gcp.hooks.bigquery import BigQueryHook from airflow.gcp.hooks.gcs import GoogleCloudStorageHook, _parse_gcs_url @@ -362,7 +364,8 @@ def get_link(self, operator, dttm): # pylint: disable=too-many-instance-attributes class BigQueryOperator(BaseOperator): """ - Executes BigQuery SQL queries in a specific BigQuery database + Executes BigQuery SQL queries in a specific BigQuery database. + This operator does not assert idempotency. :param sql: the sql code to be executed (templated) :type sql: Can receive a str representing a sql statement, @@ -746,15 +749,26 @@ def execute(self, context): conn = bq_hook.get_conn() cursor = conn.cursor() - cursor.create_empty_table( - project_id=self.project_id, - dataset_id=self.dataset_id, - table_id=self.table_id, - schema_fields=schema_fields, - time_partitioning=self.time_partitioning, - labels=self.labels, - encryption_configuration=self.encryption_configuration - ) + try: + self.log.info('Creating Table %s:%s.%s', + self.project_id, self.dataset_id, self.table_id) + cursor.create_empty_table( + project_id=self.project_id, + dataset_id=self.dataset_id, + table_id=self.table_id, + schema_fields=schema_fields, + time_partitioning=self.time_partitioning, + labels=self.labels, + encryption_configuration=self.encryption_configuration + ) + self.log.info('Table created successfully: %s:%s.%s', + self.project_id, self.dataset_id, self.table_id) + except HttpError as err: + if err.resp.status != 409: + raise + else: + self.log.info('Table %s:%s.%s already exists.', self.project_id, + self.dataset_id, self.table_id) # pylint: disable=too-many-instance-attributes @@ -917,22 +931,26 @@ def execute(self, context): conn = bq_hook.get_conn() cursor = conn.cursor() - cursor.create_external_table( - external_project_dataset_table=self.destination_project_dataset_table, - schema_fields=schema_fields, - source_uris=source_uris, - source_format=self.source_format, - compression=self.compression, - skip_leading_rows=self.skip_leading_rows, - field_delimiter=self.field_delimiter, - max_bad_records=self.max_bad_records, - quote_character=self.quote_character, - allow_quoted_newlines=self.allow_quoted_newlines, - allow_jagged_rows=self.allow_jagged_rows, - src_fmt_configs=self.src_fmt_configs, - labels=self.labels, - encryption_configuration=self.encryption_configuration - ) + try: + cursor.create_external_table( + external_project_dataset_table=self.destination_project_dataset_table, + schema_fields=schema_fields, + source_uris=source_uris, + source_format=self.source_format, + compression=self.compression, + skip_leading_rows=self.skip_leading_rows, + field_delimiter=self.field_delimiter, + max_bad_records=self.max_bad_records, + quote_character=self.quote_character, + allow_quoted_newlines=self.allow_quoted_newlines, + allow_jagged_rows=self.allow_jagged_rows, + src_fmt_configs=self.src_fmt_configs, + labels=self.labels, + encryption_configuration=self.encryption_configuration + ) + except HttpError as err: + if err.resp.status != 409: + raise class BigQueryDeleteDatasetOperator(BaseOperator): @@ -1083,11 +1101,18 @@ def execute(self, context): conn = bq_hook.get_conn() cursor = conn.cursor() - cursor.create_empty_dataset( - project_id=self.project_id, - dataset_id=self.dataset_id, - dataset_reference=self.dataset_reference, - location=self.location) + try: + self.log.info('Creating Dataset: %s in project: %s ', self.dataset_id, self.project_id) + cursor.create_empty_dataset( + project_id=self.project_id, + dataset_id=self.dataset_id, + dataset_reference=self.dataset_reference, + location=self.location) + self.log.info('Dataset created successfully.') + except HttpError as err: + if err.resp.status != 409: + raise + self.log.info('Dataset %s already exists.', self.dataset_id) class BigQueryGetDatasetOperator(BaseOperator):