This is an automated email from the ASF dual-hosted git repository.
potiuk pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new f8f53e0e94d Add ElasticsearchRemoteLogIO.from_config and register
elasticsearch scheme (#70525)
f8f53e0e94d is described below
commit f8f53e0e94dc15ace2cc4b8daa770b832b024972
Author: Yuseok Jo <[email protected]>
AuthorDate: Sat Aug 1 14:52:43 2026 +0900
Add ElasticsearchRemoteLogIO.from_config and register elasticsearch scheme
(#70525)
* Add ElasticsearchRemoteLogIO.from_config and register elasticsearch scheme
* Keep Elasticsearch host default when the config value is empty
---
providers/elasticsearch/docs/logging/index.rst | 13 +++++
providers/elasticsearch/provider.yaml | 4 ++
.../providers/elasticsearch/get_provider_info.py | 6 +++
.../providers/elasticsearch/log/es_task_handler.py | 22 ++++++++
.../unit/elasticsearch/log/test_es_task_handler.py | 62 ++++++++++++++++++++++
5 files changed, 107 insertions(+)
diff --git a/providers/elasticsearch/docs/logging/index.rst
b/providers/elasticsearch/docs/logging/index.rst
index df2a6e6be71..f3fd5486ff0 100644
--- a/providers/elasticsearch/docs/logging/index.rst
+++ b/providers/elasticsearch/docs/logging/index.rst
@@ -37,6 +37,19 @@ First, to use the handler, ``airflow.cfg`` must be
configured as follows:
[elasticsearch]
host = <host>:<port>
+On Airflow 3.x you can also route remote logging to Elasticsearch through the
provider
+dispatch mechanism by adding an ``elasticsearch://`` scheme to
+``[logging] remote_base_log_folder``:
+
+.. code-block:: ini
+
+ [logging]
+ remote_logging = True
+ remote_base_log_folder = elasticsearch://
+
+ [elasticsearch]
+ host = <host>:<port>
+
To output task logs to stdout in JSON format, the following config could be
used:
.. code-block:: ini
diff --git a/providers/elasticsearch/provider.yaml
b/providers/elasticsearch/provider.yaml
index 11de3dc9891..3e9d1933a13 100644
--- a/providers/elasticsearch/provider.yaml
+++ b/providers/elasticsearch/provider.yaml
@@ -116,6 +116,10 @@ connection-types:
logging:
-
airflow.providers.elasticsearch.log.es_task_handler.ElasticsearchTaskHandler
+remote-logging:
+ - classpath:
airflow.providers.elasticsearch.log.es_task_handler.ElasticsearchRemoteLogIO
+ scheme: elasticsearch
+
config:
elasticsearch:
description: ~
diff --git
a/providers/elasticsearch/src/airflow/providers/elasticsearch/get_provider_info.py
b/providers/elasticsearch/src/airflow/providers/elasticsearch/get_provider_info.py
index b0853a98580..fbd155a4afd 100644
---
a/providers/elasticsearch/src/airflow/providers/elasticsearch/get_provider_info.py
+++
b/providers/elasticsearch/src/airflow/providers/elasticsearch/get_provider_info.py
@@ -48,6 +48,12 @@ def get_provider_info():
}
],
"logging":
["airflow.providers.elasticsearch.log.es_task_handler.ElasticsearchTaskHandler"],
+ "remote-logging": [
+ {
+ "classpath":
"airflow.providers.elasticsearch.log.es_task_handler.ElasticsearchRemoteLogIO",
+ "scheme": "elasticsearch",
+ }
+ ],
"config": {
"elasticsearch": {
"description": None,
diff --git
a/providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_task_handler.py
b/providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_task_handler.py
index 9b4261072dc..a902bfead3a 100644
---
a/providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_task_handler.py
+++
b/providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_task_handler.py
@@ -712,6 +712,28 @@ class ElasticsearchRemoteLogIO(LoggingMixin): # noqa: D101
processors = ()
+ @classmethod
+ def from_config(cls) -> ElasticsearchRemoteLogIO:
+ """
+ Build the remote log IO from Airflow logging and ``[elasticsearch]``
configuration.
+
+ Mirrors the legacy branch in ``airflow_local_settings.py``. Unlike the
object-storage
+ backends, this does not merge ``[logging] remote_task_handler_kwargs``
IO-kwargs, matching
+ the legacy behavior for Elasticsearch.
+ """
+ return cls(
+
base_log_folder=os.path.expanduser(conf.get_mandatory_value("logging",
"base_log_folder")),
+ delete_local_copy=conf.getboolean("logging", "delete_local_logs"),
+ host=conf.get("elasticsearch", "host") or "http://localhost:9200",
+ target_index=conf.get_mandatory_value("elasticsearch",
"target_index"),
+ write_stdout=conf.getboolean("elasticsearch", "write_stdout"),
+ write_to_es=conf.getboolean("elasticsearch", "write_to_es"),
+ json_format=conf.getboolean("elasticsearch", "json_format"),
+ host_field=conf.get_mandatory_value("elasticsearch", "host_field"),
+ offset_field=conf.get_mandatory_value("elasticsearch",
"offset_field"),
+ log_id_template=conf.get_mandatory_value("elasticsearch",
"log_id_template"),
+ )
+
def __attrs_post_init__(self):
es_kwargs = get_es_kwargs_from_config()
self.client = apply_compat_with(elasticsearch.Elasticsearch(self.host,
**es_kwargs))
diff --git
a/providers/elasticsearch/tests/unit/elasticsearch/log/test_es_task_handler.py
b/providers/elasticsearch/tests/unit/elasticsearch/log/test_es_task_handler.py
index 589174c0c71..5fa6f1ad6ef 100644
---
a/providers/elasticsearch/tests/unit/elasticsearch/log/test_es_task_handler.py
+++
b/providers/elasticsearch/tests/unit/elasticsearch/log/test_es_task_handler.py
@@ -20,6 +20,7 @@ from __future__ import annotations
import dataclasses
import json
import logging
+import os
import re
from io import StringIO
from pathlib import Path
@@ -1022,3 +1023,64 @@ class TestSafeBuildStructuredLogMessage:
assert result.event == str(["a", "b"])
assert result.timestamp is not None
mock_logger.debug.assert_called_once()
+
+
+class TestElasticsearchRemoteLogIOFromConfig:
+ @conf_vars(
+ {
+ ("logging", "base_log_folder"): "~/airflow/logs",
+ ("logging", "delete_local_logs"): "True",
+ ("elasticsearch", "host"): "http://elasticsearch.example.com:9200",
+ ("elasticsearch", "target_index"): "my-logs",
+ ("elasticsearch", "write_stdout"): "True",
+ ("elasticsearch", "write_to_es"): "True",
+ ("elasticsearch", "json_format"): "True",
+ ("elasticsearch", "host_field"): "host.name",
+ ("elasticsearch", "offset_field"): "log.offset",
+ ("elasticsearch", "log_id_template"):
"{dag_id}-{task_id}-{run_id}",
+ }
+ )
+ def test_from_config(self):
+ subject = ElasticsearchRemoteLogIO.from_config()
+
+ assert subject.base_log_folder ==
Path(os.path.expanduser("~/airflow/logs"))
+ assert subject.delete_local_copy is True
+ assert subject.host == "http://elasticsearch.example.com:9200"
+ assert subject.target_index == "my-logs"
+ assert subject.write_stdout is True
+ assert subject.write_to_es is True
+ assert subject.json_format is True
+ assert subject.host_field == "host.name"
+ assert subject.offset_field == "log.offset"
+ assert subject.log_id_template == "{dag_id}-{task_id}-{run_id}"
+
+ @conf_vars(
+ {
+ ("logging", "base_log_folder"): "~/airflow/logs",
+ ("elasticsearch", "host"): "",
+ ("elasticsearch", "target_index"): "my-logs",
+ ("elasticsearch", "host_field"): "host",
+ ("elasticsearch", "offset_field"): "offset",
+ ("elasticsearch", "log_id_template"):
"{dag_id}-{task_id}-{run_id}",
+ }
+ )
+ def test_from_config_missing_host_keeps_class_default(self):
+ # An empty [elasticsearch] host must not override the class default
with "", which would
+ # make elasticsearch.Elasticsearch("") raise and silently disable
remote logging.
+ subject = ElasticsearchRemoteLogIO.from_config()
+
+ assert subject.host == "http://localhost:9200"
+
+ def test_provider_registers_elasticsearch_scheme(self):
+ from airflow.providers_manager import ProvidersManager
+
+ manager = ProvidersManager()
+ if not hasattr(manager, "remote_logging_handler_by_scheme"):
+ pytest.skip("Airflow core does not support remote logging provider
dispatch")
+
+ info = manager.remote_logging_handler_by_scheme("elasticsearch")
+
+ assert info is not None
+ assert (
+ info.classpath ==
"airflow.providers.elasticsearch.log.es_task_handler.ElasticsearchRemoteLogIO"
+ )