Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions UPDATING.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
57 changes: 14 additions & 43 deletions airflow/gcp/hooks/bigquery.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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:
"""
Expand Down
87 changes: 56 additions & 31 deletions airflow/gcp/operators/bigquery.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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):
Expand Down Expand Up @@ -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):
Expand Down