Skip to content

Add write log capacity to ElasticsearchTaskHandler & OpensearchTaskHandler #42780

Description

@Owen-CH-Leung

Description

According to the aws provider doc, when enabling remote logging, the CloudwatchTaskHandler and S3TaskHandler supports both reading & writing task logs

https://airflow.apache.org/docs/apache-airflow-providers-amazon/stable/logging/index.html

As for remote logging with Elasticsearch / Opensearch, currently only reading log is supported. Users need to deploy other softwares (such as filebeat & logstash) to ship Airflow task logs to Elasticsearch / Opensearch. Also, user would need to ensure the log messages contain a valid log_id of format {dag_id}-{task_id}-{execution_date}-{try_number} in order for reading remote log to work.

Wouldn't be nice if Airflow supports writing each task log to Elasticsearch / Opensearch, after each DAG task is completed ? Similar to S3TaskHandler, once remote logging is properly configured, DAG task log will automatically be written to, and read from Elasticsearch / Opensearch, and users need not deploy additional software to ship task logs

Use case/motivation

Similar to S3TaskHandler, ElasticsearchTaskHandler and OpensearchTaskHandler should support automatically writing task logs to destination.

Related issues

No response

Are you willing to submit a PR?

  • Yes I am willing to submit a PR!

Code of Conduct

Activity

  1. boring-cyborg commented on Oct 6, 2024

    @boring-cyborg

    Thanks for opening your first issue here! Be sure to follow the issue template! If you are willing to raise PR to address this issue please do so, no need to wait for approval.

  2. romsharon98 commented on Oct 8, 2024

    @romsharon98
    Contributor

    I would love this feature.
    When I tried to implement it myself by setting up Logstash and Filebeat as you mentioned, I encountered an issue.
    Filebeat struggled to read the log files due to the default Airflow log filename template, where the hierarchy was structured like this:

    dag_id={ ti.dag_id }/run_id={ ti.run_id }/task_id={ ti.task_id }/{%% if ti.map_index >= 0 %%}map_index={ ti.map_index }/{%% endif %%}attempt={ try_number }.log

    The hierarchy was too deep and contained too many folders (around 600,000) for Filebeat to process efficiently, causing a significant delay (hours) before the first log file was sent to Logstash. As a result, there was a major lag in log delivery.

    I believe with this structure, it’s challenging to implement effectively. However, with your proposed approach (placing all log files in a single flat directory with no depth), it would definitely work better!

    (We didn’t change the log format because we didn’t want to risk losing logs.)

  3. Owen-CH-Leung commented on Oct 8, 2024

    @Owen-CH-Leung
    ContributorAuthor

    Yeah what I propose here is - user shouldn't need to deploy additional software like filebeat, and airflow will handle shipping logs to elasticsearch for you. Just like what S3TaskHandler does

  4. romsharon98 commented on Oct 8, 2024

    @romsharon98
    Contributor

    Yeah what I propose here is - user shouldn't need to deploy additional software like filebeat, and airflow will handle shipping logs to elasticsearch for you. Just like what S3TaskHandler does

    If you will be able to do it also for the default log file template, it will be fantastic!

  5. removed
    needs-triagelabel for new issues that we didn't triage yet
    on Feb 12, 2026
  6. eladkal commented on Apr 7, 2026

    @eladkal
    Contributor

    opensearch part is completed in #64364

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions