|
23 | 23 | """ |
24 | 24 |
|
25 | 25 | import os |
| 26 | +from typing import Any, Dict |
26 | 27 |
|
27 | 28 | from airflow import models |
28 | 29 | from airflow.providers.google.cloud.operators.datastore import ( |
| 30 | + CloudDatastoreAllocateIdsOperator, CloudDatastoreBeginTransactionOperator, CloudDatastoreCommitOperator, |
29 | 31 | CloudDatastoreExportEntitiesOperator, CloudDatastoreImportEntitiesOperator, |
| 32 | + CloudDatastoreRollbackOperator, CloudDatastoreRunQueryOperator, |
30 | 33 | ) |
31 | 34 | from airflow.utils import dates |
32 | 35 |
|
|
37 | 40 | "example_gcp_datastore", |
38 | 41 | schedule_interval=None, # Override to match your needs |
39 | 42 | start_date=dates.days_ago(1), |
40 | | - tags=['example'], |
| 43 | + tags=["example"], |
41 | 44 | ) as dag: |
| 45 | + # [START how_to_export_task] |
42 | 46 | export_task = CloudDatastoreExportEntitiesOperator( |
43 | 47 | task_id="export_task", |
44 | 48 | bucket=BUCKET, |
45 | 49 | project_id=GCP_PROJECT_ID, |
46 | 50 | overwrite_existing=True, |
47 | 51 | ) |
| 52 | + # [END how_to_export_task] |
48 | 53 |
|
| 54 | + # [START how_to_import_task] |
49 | 55 | import_task = CloudDatastoreImportEntitiesOperator( |
50 | 56 | task_id="import_task", |
51 | 57 | bucket="{{ task_instance.xcom_pull('export_task')['response']['outputUrl'].split('/')[2] }}", |
52 | 58 | file="{{ '/'.join(task_instance.xcom_pull('export_task')['response']['outputUrl'].split('/')[3:]) }}", |
53 | | - project_id=GCP_PROJECT_ID |
| 59 | + project_id=GCP_PROJECT_ID, |
54 | 60 | ) |
| 61 | + # [END how_to_import_task] |
55 | 62 |
|
56 | 63 | export_task >> import_task |
| 64 | + |
| 65 | +# [START how_to_keys_def] |
| 66 | +KEYS = [ |
| 67 | + { |
| 68 | + "partitionId": {"projectId": GCP_PROJECT_ID, "namespaceId": ""}, |
| 69 | + "path": {"kind": "airflow"}, |
| 70 | + } |
| 71 | +] |
| 72 | +# [END how_to_keys_def] |
| 73 | + |
| 74 | +# [START how_to_transaction_def] |
| 75 | +TRANSACTION_OPTIONS: Dict[str, Any] = {"readWrite": {}} |
| 76 | +# [END how_to_transaction_def] |
| 77 | + |
| 78 | +# [START how_to_commit_def] |
| 79 | +COMMIT_BODY = { |
| 80 | + "mode": "TRANSACTIONAL", |
| 81 | + "mutations": [ |
| 82 | + { |
| 83 | + "insert": { |
| 84 | + "key": KEYS[0], |
| 85 | + "properties": {"string": {"stringValue": "airflow is awesome!"}}, |
| 86 | + } |
| 87 | + } |
| 88 | + ], |
| 89 | + "transaction": "{{ task_instance.xcom_pull('begin_transaction_commit') }}", |
| 90 | +} |
| 91 | +# [END how_to_commit_def] |
| 92 | + |
| 93 | +# [START how_to_query_def] |
| 94 | +QUERY = { |
| 95 | + "partitionId": {"projectId": GCP_PROJECT_ID, "namespaceId": ""}, |
| 96 | + "readOptions": { |
| 97 | + "transaction": "{{ task_instance.xcom_pull('begin_transaction_query') }}" |
| 98 | + }, |
| 99 | + "query": {}, |
| 100 | +} |
| 101 | +# [END how_to_query_def] |
| 102 | + |
| 103 | +with models.DAG( |
| 104 | + "example_gcp_datastore_operations", |
| 105 | + start_date=dates.days_ago(1), |
| 106 | + schedule_interval=None, # Override to match your needs |
| 107 | + tags=["example"], |
| 108 | +) as dag2: |
| 109 | + # [START how_to_allocate_ids] |
| 110 | + allocate_ids = CloudDatastoreAllocateIdsOperator( |
| 111 | + task_id="allocate_ids", partial_keys=KEYS, project_id=GCP_PROJECT_ID |
| 112 | + ) |
| 113 | + # [END how_to_allocate_ids] |
| 114 | + |
| 115 | + # [START how_to_begin_transaction] |
| 116 | + begin_transaction_commit = CloudDatastoreBeginTransactionOperator( |
| 117 | + task_id="begin_transaction_commit", |
| 118 | + transaction_options=TRANSACTION_OPTIONS, |
| 119 | + project_id=GCP_PROJECT_ID, |
| 120 | + ) |
| 121 | + # [END how_to_begin_transaction] |
| 122 | + |
| 123 | + # [START how_to_commit_task] |
| 124 | + commit_task = CloudDatastoreCommitOperator( |
| 125 | + task_id="commit_task", body=COMMIT_BODY, project_id=GCP_PROJECT_ID |
| 126 | + ) |
| 127 | + # [END how_to_commit_task] |
| 128 | + |
| 129 | + allocate_ids >> begin_transaction_commit >> commit_task |
| 130 | + |
| 131 | + begin_transaction_query = CloudDatastoreBeginTransactionOperator( |
| 132 | + task_id="begin_transaction_query", |
| 133 | + transaction_options=TRANSACTION_OPTIONS, |
| 134 | + project_id=GCP_PROJECT_ID, |
| 135 | + ) |
| 136 | + |
| 137 | + # [START how_to_run_query] |
| 138 | + run_query = CloudDatastoreRunQueryOperator( |
| 139 | + task_id="run_query", body=QUERY, project_id=GCP_PROJECT_ID |
| 140 | + ) |
| 141 | + # [END how_to_run_query] |
| 142 | + |
| 143 | + allocate_ids >> begin_transaction_query >> run_query |
| 144 | + |
| 145 | + begin_transaction_to_rollback = CloudDatastoreBeginTransactionOperator( |
| 146 | + task_id="begin_transaction_to_rollback", |
| 147 | + transaction_options=TRANSACTION_OPTIONS, |
| 148 | + project_id=GCP_PROJECT_ID, |
| 149 | + ) |
| 150 | + |
| 151 | + # [START how_to_rollback_transaction] |
| 152 | + rollback_transaction = CloudDatastoreRollbackOperator( |
| 153 | + task_id="rollback_transaction", |
| 154 | + transaction="{{ task_instance.xcom_pull('begin_transaction_to_rollback') }}", |
| 155 | + ) |
| 156 | + begin_transaction_to_rollback >> rollback_transaction |
| 157 | + # [END how_to_rollback_transaction] |
0 commit comments