我有一种情况,我需要在 S3 中找到一个特定的文件夹才能在 Airflow 脚本中传递给 PythonOperator。我正在使用另一个找到正确目录的 PythonOperator 来执行此操作。我可以成功地使用 xcom.push() 或 Variable.set() 并在PythonOperator 中读回它。问题是,我需要将此变量传递给使用 Python 库中代码的单独 PythonOperator。因此,我需要在 Airflow 脚本的主要部分中使用 Variable.get() 或 xcom.pull() 这个变量。我已经搜索了很多,似乎无法弄清楚这是否可能。下面是一些代码供参考:
def check_for_done_file(**kwargs):
### This function does a bunch of stuff to find the correct S3 path to
### populate target_dir, this has been verified and works
Variable.set("target_dir", done_file_list.pop())
test = Variable.get("target_dir")
print("TEST: ", test)
#### END OF METHOD, BEGIN MAIN
with my_dag:
### CALLING METHOD FROM MAIN, POPULATING VARIABLE
check_for_done_file_task = PythonOperator(
task_id = 'check_for_done_file',
python_callable …Run Code Online (Sandbox Code Playgroud) 我知道这个问题之前已经被问过,并且我已经看到了一些 SO 响应并阅读了有关该主题的 AWS 文档...我有一个 terraform 模块,它部分构建了 ECS 服务、集群、任务,和 Fargate 容器:
###############################################################################
#### EFS for added stoage
#### TODO: remove in favor of larger ephmemeral storage when terraform supports it
###############################################################################
resource "aws_efs_file_system" "test" {
creation_token = var.fargate_container_name
tags = {
Name = "test"
}
}
resource "aws_efs_access_point" "test" {
file_system_id = aws_efs_file_system.test.id
root_directory {
path = "/"
}
}
resource "aws_efs_mount_target" "test" {
count = 3
file_system_id = aws_efs_file_system.test.id
subnet_id = local.directory_subnet_ids[count.index]
security_groups = [aws_security_group.test_ecs.id]
}
###############################################################################
#### ECS …Run Code Online (Sandbox Code Playgroud) 我正在尝试在 AWS EC2 实例上安装气流。网络上的各种来源似乎都很好地记录了该过程,但是,我在“pip install”气流后遇到了问题;执行命令“airflow initdb”时出现以下错误:
[2019-09-25 13:22:02,329] {__init__.py:51} INFO - Using executor SequentialExecutor
Traceback (most recent call last):
File "/home/cloud-user/.local/bin/airflow", line 22, in <module>
from airflow.bin.cli import CLIFactory
File "/home/cloud-user/.local/lib/python2.7/site-packages/airflow/bin/cli.py", line 68, in <module>
from airflow.www_rbac.app import cached_app as cached_app_rbac
File "/home/cloud-user/.local/lib/python2.7/site-packages/airflow/www_rbac/app.py", line 26, in <module>
from flask_appbuilder import AppBuilder, SQLA
File "/home/cloud-user/.local/lib/python2.7/site-packages/flask_appbuilder/__init__.py", line 5, in <module>
from .base import AppBuilder
File "/home/cloud-user/.local/lib/python2.7/site-packages/flask_appbuilder/base.py", line 5, in <module>
from .api.manager import OpenApiManager
File "/home/cloud-user/.local/lib/python2.7/site-packages/flask_appbuilder/api/__init__.py", line 11, in <module>
from marshmallow_sqlalchemy.fields …Run Code Online (Sandbox Code Playgroud) 我正在尝试从触发的 Airflow SimpleHttpOperator 接收 HTTP 响应代码。我已经看到使用 'lambda' 类型的示例,并且正在通过查看响应正文来这样做,但我希望能够将响应代码传递给一个函数。我当前的代码(其中 90% 来自 example_http_operator):
import json
from datetime import timedelta
from airflow import DAG
from airflow.operators.http_operator import SimpleHttpOperator
from airflow.sensors.http_sensor import HttpSensor
from airflow.utils.dates import days_ago
default_args = {
'owner': 'airflow',
'depends_on_past': False,
'start_date': days_ago(2),
'email': ['me@company.com'],
'email_on_failure': False,
'email_on_retry': False,
'retries': 0,
}
dag = DAG(dag_id='kick_off_java_task', default_args=default_args)
kickoff_task = SimpleHttpOperator(
task_id='kick_off_c2c_java_task',
http_conn_id='test',
method='GET',
endpoint='',
data={ "command": "run" },
response_check=lambda response: True if "Ok Message" in response.text else False,
headers={},
xcom_push=False, …Run Code Online (Sandbox Code Playgroud)