Проблема, с которой я столкнулся, заключается в наличии нескольких баз данных и связанных с ними таблиц в каждой базе данных. Например, у меня есть список баз данных: my_sql_bds = ["назначения", "карты", "клинические формы", "онбординг", "производство", "системный менеджмент"]
В в каждом из них есть несколько таблиц.
Проблема, которую я пытаюсь решить, заключается в получении Information_schema для каждой таблицы в каждой базе данных.
Сейчас я сначала удаляю все файлы предыдущего дня:
Код: Выделить всё
def delete_temp_gcs_objects(gcs_hook, bucket_name, objects, delimiter):
try:
objects = gcs_hook.list(bucket_name, prefix=objects, delimiter=delimiter)
for source_object in objects:
gcs_hook.delete(bucket_name, source_object)
except Exception as e:
print(e)
Код: Выделить всё
def setup_connection(MYSQL_address: str) -> None:
connection = Connection(
conn_id=CONNECTION_ID,
host=ip_address,
login=DB_USER_NAME,
password=DB_USER_PASSWORD,
port=DB_PORT,
)
"""Make sure I kill connection after all jobs are done."""
session = Session()
log.info("Removing connection %s if it exists", CONNECTION_ID)
query = session.query(Connection).filter(Connection.conn_id == CONNECTION_ID)
query.delete()
session.add(connection)
session.commit()
log.info("Connection %s created", CONNECTION_ID)
setup_connection_task = setup_connection(pull_mysql_address)
Код: Выделить всё
for db, table in my_sql_bds.items():
MySQLToGCSOperator(
task_id = "my_sql_to_gcs",
mysql_conn_id = CONNECTION_ID,
sql= SCHEMA_EXTRACTION
bucket= DESTINATION_BUCKET,
filename= f"/temp/mysql_transfer/{db}_information_schema.json"
export_format="json"
)
Код: Выделить всё
tables_dict = {
"appointments":[ "call_schedule", "call_config"],
"onboarding": ["address", "address_history"]
}
Код: Выделить всё
for db, table in tables_dict.items():
for i in table:
MySQLToGCSOperator(
task_id = "my_sql_to_gcs",
mysql_conn_id = CONNECTION_ID,
sql= COLUMN_EXTRACTION
bucket= DESTINATION_BUCKET,
filename= f"/temp/mysql_transfer/{db}/{i}/{db}_{i}_information_schema.json"
export_format="json"
Изначально в GCS была всего одна таблица, поэтому мне был нужен только этот код.
Код: Выделить всё
ops_mysql_to_gcs = MySQLToGCSOperator(
task_id="ops_mysql_to_gcs",
mysql_conn_id=CONNECTION_ID,
sql=OPERATIONAL_STAGING_EXTRACTION,
bucket=BUCKET_NAME,
filename=STAGING_FILENAME,
export_format="json",
)
Подробнее здесь: https://stackoverflow.com/questions/787 ... end-to-gcs