Google Cloud Datastore Operators¶
Firestore in Datastore mode is a NoSQL document database built for automatic scaling, high performance, and ease of application development.
For more information about the service visit Datastore product documentation
Prerequisite Tasks¶
To use these operators, you must do a few things:
Select or create a Cloud Platform project using the Cloud Console.
Enable billing for your project, as described in the Google Cloud documentation.
Enable the API, as described in the Cloud Console documentation.
Install API libraries via pip.
pip install 'apache-airflow[google]'Detailed information is available for Installation.
Export Entities¶
To export entities from Google Cloud Datastore to Cloud Storage use
CloudDatastoreExportEntitiesOperator
export_task = CloudDatastoreExportEntitiesOperator(
task_id="export_task",
bucket=BUCKET_NAME,
project_id=PROJECT_ID,
overwrite_existing=True,
)
Import Entities¶
To import entities from Cloud Storage to Google Cloud Datastore use
CloudDatastoreImportEntitiesOperator
import_task = CloudDatastoreImportEntitiesOperator(
task_id="import_task",
bucket="{{ task_instance.xcom_pull('export_task')['response']['outputUrl'].split('/')[2] }}",
file="{{ '/'.join(task_instance.xcom_pull('export_task')['response']['outputUrl'].split('/')[3:]) }}",
project_id=PROJECT_ID,
)
Allocate Ids¶
To allocate IDs for incomplete keys use
CloudDatastoreAllocateIdsOperator
allocate_ids = CloudDatastoreAllocateIdsOperator(
task_id="allocate_ids", partial_keys=KEYS, project_id=PROJECT_ID
)
An example of a partial keys required by the operator:
KEYS = [
{
"partitionId": {"projectId": PROJECT_ID, "namespaceId": ""},
"path": {"kind": "airflow"},
}
]
Begin transaction¶
To begin a new transaction use
CloudDatastoreBeginTransactionOperator
TRANSACTION_OPTIONS = {"readWrite": {}}
begin_transaction = CloudDatastoreBeginTransactionOperator(
task_id="begin_transaction",
transaction_options=TRANSACTION_OPTIONS,
project_id=PROJECT_ID,
)
Warning
Datastore transactions expire after 270 seconds or after 60 seconds of inactivity.
Airflow does not guarantee that a downstream task will start before those limits.
Do not pass a transaction handle between tasks. Run all operations belonging to a
transaction within one task using
DatastoreHook or a Datastore client.
See Datastore transaction limits.
Commit transaction¶
To commit a transaction, optionally creating, deleting or modifying some entities
use CloudDatastoreCommitOperator
commit_task = CloudDatastoreCommitOperator(task_id="commit_task", body=COMMIT_BODY, project_id=PROJECT_ID)
An example of a commit information required by the operator:
COMMIT_BODY = {
"mode": "TRANSACTIONAL",
"mutations": [
{
"insert": {
"key": KEYS[0],
"properties": {"string": {"stringValue": "airflow is awesome!"}},
}
}
],
"singleUseTransaction": {"readWrite": {}},
}
Run query¶
To run a query for entities use
CloudDatastoreRunQueryOperator
run_query = CloudDatastoreRunQueryOperator(task_id="run_query", body=QUERY, project_id=PROJECT_ID)
An example of a query required by the operator:
QUERY = {
"partitionId": {"projectId": PROJECT_ID, "namespaceId": "query"},
"query": {},
}
Roll back transaction¶
To roll back a transaction
use CloudDatastoreRollbackOperator
This low-level operator requires a transaction handle that is still active when task execution begins. Given such a handle, configure the operator as follows:
rollback_transaction = CloudDatastoreRollbackOperator(
task_id="rollback_transaction",
transaction=TRANSACTION_ID,
project_id=PROJECT_ID,
)
Get operation state¶
To get the current state of a long-running operation use
CloudDatastoreGetOperationOperator
get_operation = CloudDatastoreGetOperationOperator(
task_id="get_operation", name="{{ task_instance.xcom_pull('export_task')['name'] }}"
)
Delete operation¶
To delete an operation use
CloudDatastoreDeleteOperationOperator
delete_export_operation = CloudDatastoreDeleteOperationOperator(
task_id="delete_export_operation",
name="{{ task_instance.xcom_pull('export_task')['name'] }}",
trigger_rule=TriggerRule.ALL_DONE,
)
References¶
For further information, take a look at: