【问题标题】:Unable to import custom operators from plugins folder airflow无法从插件文件夹气流导入自定义运算符
【发布时间】:2021-05-24 10:13:34
【问题描述】:

我是气流新手,并试图从我的 dag 的插件文件夹中导入自定义运算符。 下面是文件结构:

├── dags
│   ├── my_dag.py
├── myrequirements.txt
├── plugins
│   ├── __init__.py
│   ├── my_airflow_plugin.py
│   └── operators
│       ├── __int__.py
│       └── my_airflow_operator.py

my_dag.py

from operators.my_airflow_operator import AwsLambdaInvokeFunctionOperator

my_airflow_plugin.py

from airflow.plugins_manager import AirflowPlugin
from operators.my_airflow_operator import AwsLambdaInvokeFunctionOperator

class lambda_operator(LambdaOperator):
    pass
                    
class my_plugin(AirflowPlugin):
                    
    name = 'my_airflow_plugin'
    operators = [lambda_operator]

my_airflow_operator.py

from airflow.models import BaseOperator
from airflow.utils.decorators import apply_defaults
from airflow.contrib.hooks.aws_lambda_hook import AwsLambdaHook


class AwsLambdaExecutionError(Exception):
    """
    Raised when there is an error executing the function.
    """


class AwsLambdaPayloadError(Exception):
    """
    Raised when there is an error with the Payload object in the response.
    """


class AwsLambdaInvokeFunctionOperator(BaseOperator):
    """
    Invoke AWS Lambda functions with a JSON payload.
    The check_success_function signature should be a single param which will receive a dict.
    The dict will be the "Response Structure" described in
    """
    
    def succeeded(response):
        payload = json.loads(response['Payload'].read())
        # do something with payload
    
    @apply_defaults
    def __init__(
        self,
        function_name,
        region_name,
        payload,
        check_success_function,
        log_type="None",
        qualifier="$LATEST",
        aws_conn_id=None,
        *args,
        **kwargs,
    ):
        super().__init__(*args, **kwargs)
        self.function_name = function_name
        self.region_name = region_name
        self.payload = payload
        self.log_type = log_type
        self.qualifier = qualifier
        self.check_success_function = check_success_function
        self.aws_conn_id = aws_conn_id

    def get_hook(self):
        """
        Initialises an AWS Lambda hook
        :return: airflow.contrib.hooks.AwsLambdaHook
        """
        return AwsLambdaHook(
            self.function_name,
            self.region_name,
            self.log_type,
            self.qualifier,
            aws_conn_id=self.aws_conn_id,
        )

    def execute(self, context):
        self.log.info("AWS Lambda: invoking %s", self.function_name)

        response = self.get_hook().invoke_lambda(self.payload)

        try:
            self._validate_lambda_api_response(response)
            self._validate_lambda_response_payload(response)
        except (AwsLambdaExecutionError, AwsLambdaPayloadError) as e:
            self.log.error(response)
            raise e

        self.log.info("AWS Lambda: %s succeeded!", self.function_name)



    def _validate_lambda_response_payload(self, response):
        """
        Call a user provided function to validate the Payload object for errors.
        :param response: HTTP Response from AWS Lambda.
        :type response: dict
        :return: None
        """
        if not self.check_success_function(response):
            raise AwsLambdaPayloadError(
                "AWS Lambda: error validating response payload!"
            )

但我收到此错误: 没有名为“操作员”的模块

我尝试将 my_dag.py 中的导入语句更改为:

from airflow.operators.my_airflow_plugin import AwsLambdaInvokeFunctionOperator

我收到此错误 没有名为“airflow.operators.my_airflow_plugin”的模块

有人可以建议这里有什么问题吗?(气流版本是 1.10.12)

init.py 文件为空

【问题讨论】:

  • 你运行的是什么 Airflow 版本?
  • 请提供原始错误的完整追溯。看起来与 my_airflow_plugin.py 的第二行有关。还包括plugins/ __init__.py的内容
  • AIrflow 版本为 1.10.12。并且 plugins/ init.py 是空的。错误仅此而已。没有名为“操作员”的模块

标签: python plugins airflow


【解决方案1】:

你需要这个插件做什么?是否只是在 DAG 中公开您的自定义运算符?对于 Airflow 架构中的插件,这并不是真正的 the purpose,从版本 2 开始不再支持。

您可以使用运算符创建模块,然后使用该模块将其导入 DAG,而不是使用插件。你可以关注this instructions 来做。

【讨论】:

  • 嗨!谢谢。是的,我正在使用插件在我的 DAG 中公开自定义运算符。现在我将 my_airflow_operator.py(custom operator) 从 plugins 移动到 dags/ 并将我的 DAG 中的 import 语句更改为 from my_airflow_operator import AwsLambdaInvokeFunctionOperator 这有效。谢谢!
【解决方案2】:

我的气流版本是 1.10.12,但仍然从插件文件夹访问自定义运算符显示“未找到模块错误”。最后,当我将自定义操作符文件移动到 dags 文件夹时,它起作用了。

【讨论】:

  • 我也有同样的问题。你能告诉我你在 DAGs 文件夹中的结构以及如何将这些插件连接到 DAG 吗?谢谢!
猜你喜欢
  • 2023-02-15
  • 2018-12-24
  • 2020-09-26
  • 1970-01-01
  • 2021-08-19
  • 2023-03-15
  • 2021-10-13
  • 2022-07-06
  • 1970-01-01
相关资源
最近更新 更多