我试图找到一种方法,连接池管理外部连接创建的气流。
气流版本: 2.1.0
Python版本: 3.9.5
气流数据库: SQLite
创建外部连接: MySQL和雪花
我知道airflow.cfg文件中有一些属性
sql_alchemy_pool_enabled = True
sql_alchemy_pool_size = 5但是这些属性是用来管理气流内部DB的,在我的例子中,这是SQLite。
我有很少的任务是读取或写入MySQL和雪花的数据。
snowflake_insert = SnowflakeOperator(
task_id='insert_snowflake',
dag=dag,
snowflake_conn_id=SNOWFLAKE_CONN_ID,
sql="Some Insert query",
warehouse=SNOWFLAKE_WAREHOUSE,
database=SNOWFLAKE_DATABASE,
schema=SNOWFLAKE_SCHEMA,
role=SNOWFLAKE_ROLE
)和
insert_mysql_task = MySqlOperator(task_id='insert_record', mysql_conn_id='mysql_default', sql="some insert query", dag=dag)从MySQL读取数据
def get_records():
mysql_hook = MySqlHook(mysql_conn_id="mysql_default")
records = mysql_hook.get_records(sql=r"""Some select query""")
print(records)我观察到的是为每个任务创建了一个新的会话(在同一个进程中有多个任务),对于MySQL还没有验证相同的任务。
是否有一种方法来维护外部连接的连接池(在我的例子中是雪花和MySQL),或者有任何其他方式在同一个会话中运行同一DAG中的所有查询?
谢谢
发布于 2021-06-14 10:48:43
气流提供使用水池作为将并发限制为外部服务的方法。
您可以通过UI:->管理->池创建一个池
或者使用CLI:
airflow pools set NAME slots池有槽,这些槽定义使用资源的任务可以并行运行的任务数。如果池已满,则任务将排队,直到打开一个插槽。
要在运算符中使用池,只需将pool=Name添加到操作符中即可。
在您的例子中,假设Pool是以雪花的名称创建的,那么:
snowflake_insert = SnowflakeOperator(
task_id='insert_snowflake',
dag=dag,
snowflake_conn_id=SNOWFLAKE_CONN_ID,
sql="Some Insert query",
warehouse=SNOWFLAKE_WAREHOUSE,
database=SNOWFLAKE_DATABASE,
schema=SNOWFLAKE_SCHEMA,
role=SNOWFLAKE_ROLE,
pool='snowflake',
)注意,默认情况下,任务占用池中的一个槽,但这是可配置的。如果使用pool_slots示例,任务可能占用超过一个时隙:
snowflake_insert = SnowflakeOperator(
task_id='insert_snowflake',
...
pool='snowflake',
pool_slots=2,
)https://stackoverflow.com/questions/67968288
复制相似问题