Как просмотреть несколько таблиц для отправки в GCS?Python

Программы на Python
Anonymous
Как просмотреть несколько таблиц для отправки в GCS?

Сообщение Anonymous »

Сейчас у меня есть задача, по которой мне нужно переместить внешние таблицы в GCS. Для этого я использую MySQLToGCSOperator, чтобы переместить мои внешние таблицы и схему в GCS; такие вещи, как таблицы, столбцы, ПК, ФК. Чтобы выполнить этот процесс, я собираюсь использовать dag для автоматизации этого процесса.
Проблема, с которой я столкнулся, заключается в наличии нескольких баз данных и связанных с ними таблиц в каждой базе данных. Например, у меня есть список баз данных: 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)
Оттуда я настраиваю внешнее соединение с MySQL:

Код: Выделить всё

    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)
Я хочу создать список словарей для каждой базы данных, а затем выполнить вызов mysql to gcs, чтобы поместить данные gcs для схемы каждой таблицы.

Код: Выделить всё

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"]
}
Тогда мне придется сделать еще один вызов с помощью MySQLToGCSOperator, чтобы получить столбцы:

Код: Выделить всё

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

Вернуться в «Python»