自定义Airflow Operator重写__new__方法参数为空、__init__无日志问题咨询
问题解答
1. 能否让重写的__new__方法正常接收参数?
可以实现,你需要在自定义Operator类中实现__getnewargs_ex__(Python 3.3+支持)魔法方法,指定pickle序列化时需要保存的参数,反序列化时这些参数会自动传递给__new__方法,示例代码如下:
class TestOperator(BaseOperator): def __new__(cls, *args, **kwargs): logger.info("in new") logger.info(f"kwargs: {kwargs}") logger.info(f"args: {args}") instance = super().__new__(cls) # 如果需要触发__init__调用可以手动加这行 # instance.__init__(*args, **kwargs) return instance # 新增该方法指定序列化传递给__new__的参数 def __getnewargs_ex__(self): # 第一个返回值是位置参数元组,第二个是关键字参数字典 return (), {"dag": self.dag, "task_id": self.task_id} def __init__(self, *args, **kwargs): logger.info("in init") super().__init__(*args, **kwargs) def execute(self, context): logger.info("in execute")
2. 无法获取参数的底层原因
你在运行日志中看到的__new__调用,并非DAG定义时创建TestOperator实例的那一次调用,而是Airflow Worker节点反序列化序列化对象的过程:
- Airflow的Scheduler在解析DAG文件时,会执行你写的
t1 = TestOperator(...)代码,此时__new__和__init__都会正常接收参数并执行,这部分日志默认不会输出到TaskInstance的运行日志中,你没有采集到而已。 - Python标准的
pickle序列化模块反序列化对象的默认逻辑是:直接调用类的__new__方法创建空实例,不传递任何参数,不会自动调用__init__方法,而是直接把序列化时存储的实例属性赋值给新实例,这就是你看到__new__参数为空的原因。
3. __init__没有日志输出的原因
和上面的序列化逻辑一致,pickle反序列化过程默认不会调用类的__init__方法,所以对应的日志自然不会输出,并非你的代码逻辑错误。
需求替代方案建议
你需要动态调整继承的父Operator的场景,并不推荐用重写__new__的方式实现,Airflow Operator的序列化/反序列化逻辑会增加实现复杂度,更建议用以下方案替代:
- 用工厂方法模式实现,定义一个返回不同Operator实例的工厂函数,在DAG定义阶段就根据条件生成对应类型的Operator
- 用组合替代继承,在自定义Operator内部持有对应父Operator的实例,在
execute方法中调用对应实例的逻辑,运行时可以动态切换内部实例
内容的提问来源于stack exchange,提问作者FloG
相关产品推荐
相关产品推荐

