2
votes

In my airflow dag, I have an ecs_operator task followed by python operator task. I want to push some messages from ECS task to python task using xcom feature of airflow. I tried the option do_xcom_push=True with no result. Find below sample dag.

dag = DAG(
    dag_name, default_args=default_args, schedule_interval=None)
start = DummyOperator(task_id = 'start'
                   ,dag =dag)
end = DummyOperator(task_id = 'end'
                   ,dag =dag)
ecs_operator_args = {
    'launch_type': 'FARGATE',
    'task_definition': 'task-def:2',
    'cluster': 'cluster-name',
    'region_name': 'region',
    'network_configuration': {
        'awsvpcConfiguration':
            {}
    }
}
ecs_task = ECSOperator(
    task_id='x_com_test'
    ,**ecs_operator_args
    ,do_xcom_push=True
    ,params={'my_param': 'Parameter-1'}
    ,dag=dag)


def pull_function(**kwargs):
    ti = kwargs['ti']
    msg = ti.xcom_pull(task_ids='x_com_test',key='the_message')
    print("received message: '%s'" % msg)

pull_task = PythonOperator(
    task_id='pull_task',
    python_callable=pull_function,
    provide_context=True,
    dag=dag)

start >> ecs_task >> pull_task >> end
2

2 Answers

1
votes

You need to setup a cloudwatch log group for the container.

ECSOperator needs to be extended to support pushing to xcom:

from collections import deque
from airflow.utils import apply_defaults
from airflow.contrib.operators.ecs_operator import ECSOperator


class MyECSOperator(ECSOperator):
    @apply_defaults
    def __init__(self, xcom_push=False, **kwargs):
        super(CLECSOperator, self).__init__(**kwargs)
        self.xcom_push_flag = xcom_push

    def execute(self, context):
        super().execute(context)
        if self.xcom_push_flag:
            return self._last_log_event()

    def _last_log_event(self):
        if self.awslogs_group and self.awslogs_stream_prefix:
            task_id = self.arn.split("/")[-1]
            stream_name = "{}/{}".format(self.awslogs_stream_prefix, task_id)
            events = self.get_logs_hook().get_log_events(self.awslogs_group, stream_name)
            last_event = deque(events, maxlen=1).pop()
            return last_event["message"]


dag = DAG(
    dag_name, default_args=default_args, schedule_interval=None)
start = DummyOperator(task_id = 'start'
                   ,dag =dag)
end = DummyOperator(task_id = 'end'
                   ,dag =dag)
ecs_operator_args = {
    'launch_type': 'FARGATE',
    'task_definition': 'task-def:2',
    'cluster': 'cluster-name',
    'region_name': 'region',
    'awslogs_group': '/aws/ecs/myLogGroup',
    'awslogs_stream_prefix': 'myStreamPrefix',
    'network_configuration': {
        'awsvpcConfiguration':
            {}
    }
}
ecs_task = MyECSOperator(
    task_id='x_com_test'
    ,**ecs_operator_args
    ,xcom_push=True
    ,params={'my_param': 'Parameter-1'}
    ,dag=dag)


def pull_function(**kwargs):
    ti = kwargs['ti']
    msg = ti.xcom_pull(task_ids='x_com_test',key='return_value')
    print("received message: '%s'" % msg)

pull_task = PythonOperator(
    task_id='pull_task',
    python_callable=pull_function,
    provide_context=True,
    dag=dag)

start >> ecs_task >> pull_task >> end

ecs_task will take the last event from the log group before finishing executing, and push it to xcom.

1
votes

Apache-AWS has a new commit that pretty much implements what @Бојан-Аџиевски mentioned above, so you don't need to write your custom ECSOperator. Available as of version 1.1.0

All you gotta do is to provide the do_xcom_push=True when calling the ECSOperator and provide the correct awslogs_group and awslogs_stream_prefix.

Make sure your awslogs_stream_prefix follows the following format:

prefix-name/container-name

As this is what ECS directs logs to.