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"
+        )

Reply via email to