【发布时间】:2020-09-26 19:47:24
【问题描述】:
我不是专业的 Python 开发人员,所以我只是概述我的步骤。
设置部分
我为 apache-airflow 创建了一个目录 ~/Desktop/airflow 并制作了
导出 AIRFLOW_HOME=~/Desktop/airflow
然后我创建了 venv 使用
python3 -m venv ~/Desktop/airflow
结果是
然后我做了
源 bin/激活
pip3 安装 apache-airflow==1.10.9
气流初始化数据库
结果是
在我的 airflow.cfg 文件中,我检查了 dags 和插件目录。我在 $AIRFLOW_HOME/Desktop/airflow
中创建了 dags 和插件目录我启动了 airflow webserver 和 scheduler 并确保一切正常。
自定义插件部分
我发现了很多创建气流插件的方法。我尝试了所有可能的方法。让我们开始吧。
第一个是在(first_plugin)项目中创建一个插件文件夹,然后创建一个python文件(first_operator.py)
import logging
from airflow.operators import BaseOperator
from airflow.utils.decorators import apply_defaults
from airflow.plugins_manager import AirflowPlugin
log = logging.getLogger(__name__)
class FirstOperator(BaseOperator):
@apply_defaults
def __init__(self, *args, **kwargs):
super(FirstOperator, self).__init__(*args, **kwargs)
def execute(self, context):
log.info("Hello World!")
class FirstOperatorPlugin(AirflowPlugin):
name = "first_plugin"
operators = [FirstOperator]
看起来像
然后我只需将我的插件文件夹 (first_plugin) 移动到 $AIRFLOW_HOME/DESKTOP/airflow/plugins 并重新启动气流网络服务器和调度程序。
现在是时候使用我的自定义运算符创建自定义 dag。如何正确导入插件是一个挑战。有很多方法可以导入自定义运算符。我会展示我的尝试。
- from airflow.operators import FirstOperator - 已弃用
- 从airflow.operators.first_plugin导入FirstOperator
- 从airflow.operators.first_operator导入FirstOperator
- 从 first_plugin.first_operator 导入 FirstOperator
在 Pycharm IDE 中导入时,这些方法都没有帮助我。例如,
从airflow.operators.first_plugin导入FirstOperator
但我确定如果我忽略导入行并将我的自定义 dag 放入 dags 文件夹中,它会正常工作。 (我试过了)。此外,我决定检查气流日志(在调试模式下)。
重启气流网络服务器时看到的日志
我花了 2 天,我仍然没有任何解决方案。可能你们告诉我尝试其他方法。我试过了。
https://pybit.es/introduction-airflow.html
它们都是可行的方法,但没有一个能解决我的 IDE 导入问题。
【问题讨论】: