Airflow XCOM传递Oracle Datetime字段时触发AttributeError报错求助
问题详情
我是Airflow新手,正在尝试使用OracleOperator从Oracle数据库查询数据。我的代码如下:
from datetime import datetime, timedelta from airflow import DAG from airflow.providers.oracle.operators.oracle import OracleOperator default_args = { 'owner': 'owner_name', 'retries': 5, 'retry_delay': timedelta(minutes=10) } with DAG( 'dag_name', default_args=default_args, schedule_interval='0 1 * * *', start_date=datetime(2023, 4, 10), ) as dag: task = OracleOperator( task_id='task_name', oracle_conn_id='connection_name', sql=""" select some_datetime_field, some_data_field from some_table """ ) task
但运行时出现如下报错:
Traceback (most recent call last): File "/home/airflow/.local/lib/python3.7/site-packages/airflow/utils/session.py", line 72, in wrapper return func(*args, **kwargs) File "/home/airflow/.local/lib/python3.7/site-packages/airflow/models/taskinstance.py", line 2305, in xcom_push session=session, File "/home/airflow/.local/lib/python3.7/site-packages/airflow/utils/session.py", line 72, in wrapper return func(*args, **kwargs) File "/home/airflow/.local/lib/python3.7/site-packages/airflow/models/xcom.py", line 240, in set map_index=map_index, File "/home/airflow/.local/lib/python3.7/site-packages/airflow/models/xcom.py", line 627, in serialize_value return json.dumps(value, cls=XComEncoder).encode("UTF-8") File "/usr/local/lib/python3.7/json/__init__.py", line 238, in dumps **kw).encode(obj) File "/home/airflow/.local/lib/python3.7/site-packages/airflow/utils/json.py", line 176, in encode return super().encode(o) File "/usr/local/lib/python3.7/json/encoder.py", line 199, in encode chunks = self.iterencode(o, _one_shot=True) File "/usr/local/lib/python3.7/json/encoder.py", line 257, in iterencode return _iterencode(o, 0) File "/home/airflow/.local/lib/python3.7/site-packages/airflow/utils/json.py", line 153, in default CLASSNAME: o.__module__ + "." + o.__class__.__qualname__, AttributeError: 'datetime.datetime' object has no attribute '__module__'
删除SQL语句中的some_datetime_field字段后任务即可正常运行。想了解报错原因,以及如何通过XCOM将SQL查询得到的datetime字段传递给其他任务?
报错原因
OracleOperator默认会把查询结果推送到XCOM,而Airflow的XComEncoder在序列化datetime对象时出现了问题——这里的datetime对象不是Python标准库的datetime.datetime(标准库对象自带__module__属性),而是Oracle驱动返回的自定义datetime类型,该类型没有__module__属性,导致序列化失败。
解决方案
有两种实用方法可以解决这个问题:
方法一:SQL层转换datetime为字符串
直接在查询语句里将datetime字段格式化为字符串,返回的字符串类型可被XCOM正常序列化:
sql=""" select to_char(some_datetime_field, 'YYYY-MM-DD HH24:MI:SS') as some_datetime_str, some_data_field from some_table """
后续任务获取XCOM数据后,再将字符串转回datetime类型即可。
方法二:自定义Operator处理查询结果
继承OracleOperator,重写execute方法,手动处理datetime字段的序列化逻辑:
from datetime import datetime from airflow.providers.oracle.operators.oracle import OracleOperator class CustomOracleOperator(OracleOperator): def execute(self, context): result = super().execute(context) # 遍历结果处理datetime对象 processed_result = [] for row in result: processed_row = [] for item in row: if isinstance(item, datetime): processed_row.append(item.isoformat()) else: processed_row.append(item) processed_result.append(tuple(processed_row)) # 将处理后的结果推送到XCOM context['ti'].xcom_push(key='processed_result', value=processed_result) return processed_result
之后在DAG中使用这个自定义Operator替代原OracleOperator即可。
内容的提问来源于stack exchange,提问作者Phudit Chalekarn
相关产品推荐
相关产品推荐

