mongoimport JSON do Google Cloud Storage em uma tarefa do Airflow
Parece que mover dados do GCS para o MongoDB não é comum, já que não há muita documentação sobre isso. Temos a seguinte tarefa que passamos como o python_callablepara um operador Python - essa tarefa move os dados do BigQuery para o GCS como JSON:
def transfer_gcs_to_mongodb(table_name):
# connect
client = bigquery.Client()
bucket_name = "our-gcs-bucket"
project_id = "ourproject"
dataset_id = "ourdataset"
destination_uri = f'gs://{bucket_name}/{table_name}.json'
dataset_ref = bigquery.DatasetReference(project_id, dataset_id)
table_ref = dataset_ref.table(table_name)
configuration = bigquery.job.ExtractJobConfig()
configuration.destination_format = 'NEWLINE_DELIMITED_JSON'
extract_job = client.extract_table(
table_ref,
destination_uri,
job_config=configuration,
location="US",
) # API request
extract_job.result() # Waits for job to complete.
print("Exported {}:{}.{} to {}".format(project_id, dataset_id, table_name, destination_uri))
Esta tarefa está obtendo dados com sucesso no GCS. No entanto, estamos travados agora quando se trata de como executar mongoimportcorretamente, para colocar esses dados no MongoDB. Em particular, parece que mongoimportnão é possível apontar para o arquivo no GCS, mas, em vez disso, ele deve ser baixado localmente primeiro e, em seguida, importado para o MongoDB.
Como isso deve ser feito no Airflow? Devemos escrever um script de shell que baixa o JSON do GCS e, em seguida, executa mongoimportcom as urisinalizações corretas e todas as corretas? Ou existe outra maneira de executar mongoimportno Airflow que está faltando?
Respostas
Você não precisa escrever um script de shell para fazer o download do GCS. Você pode simplesmente usar o GCSToLocalFilesystemOperator e então abrir o arquivo e gravá-lo no mongo usando a função insert_many do MongoHook.
Eu não testei, mas deveria ser algo como:
mongo = MongoHook(conn_id=mongo_conn_id)
with open('file.json') as f:
file_data = json.load(f)
mongo.insert_many(file_data)
Isso é para um canal de: BigQuery -> GCS -> Sistema de arquivos local -> MongoDB.
Você também pode fazer isso na memória como: BigQuery -> GCS -> MongoDB se preferir.