Airflow可延迟传感器出现AttributeError: __aenter__问题求助
错误原因分析
- 你在异步
run方法中调用了同步实现的MsSqlIntegratedHook及其游标:该Hook和游标仅支持同步上下文管理,没有定义异步上下文所需的async def __aenter__和async def __aexit__方法,因此用async with会触发__aenter__属性错误;而普通with在异步函数中执行同步阻塞IO操作,会干扰Airflow异步事件循环,同时引发上下文适配问题。 - 你手动添加的
__aenter__和__aexit__是同步魔法方法,并非异步版本,无法适配异步上下文环境。
解决方案
方案1:使用异步MS SQL Hook(推荐)
如果你的Airflow环境存在支持异步的MS SQL Hook(如第三方异步实现或Airflow官方后续版本的异步Hook),直接替换为异步Hook即可,示例代码:
async def run(self): while True: # 假设存在AsyncMsSqlIntegratedHook异步实现 async with AsyncMsSqlIntegratedHook(mssql_conn_id=self.mssql_conn_id).get_async_cursor() as cursor: await cursor.execute(self.sql) rows = await cursor.fetchone() if rows[0] > 0: yield TriggerEvent(True) await asyncio.sleep(self.sleep_interval)
方案2:将同步操作包装到线程池(兼容现有同步Hook)
把同步的Hook查询操作放到asyncio线程池中执行,避免阻塞异步事件循环,示例代码:
import asyncio async def run(self): while True: # 定义同步查询逻辑 def check_condition(): with MsSqlIntegratedHook(mssql_conn_id=self.mssql_conn_id).get_cursor() as cursor: cursor.execute(self.sql) rows = cursor.fetchone() return rows[0] > 0 # 在线程池中调度同步操作 condition_met = await asyncio.get_event_loop().run_in_executor(None, check_condition) if condition_met: yield TriggerEvent(True) await asyncio.sleep(self.sleep_interval)
关键注意事项
- 禁止在异步函数中直接执行同步阻塞IO操作(如数据库查询),必须通过线程池/进程池调度,否则会卡住整个Airflow异步调度器。
- 若自定义异步Hook,必须实现异步上下文管理方法:
class AsyncMsSqlIntegratedHook(BaseHook): async def __aenter__(self): # 实现异步初始化连接逻辑 self.connection = await self._get_async_connection() return self async def __aexit__(self, exc_type, exc_val, exc_tb): # 实现异步关闭连接逻辑 await self.connection.close()
内容的提问来源于stack exchange,提问作者Saugat Mukherjee
相关产品推荐
相关产品推荐

