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 a472c83d715 Add WasbRemoteLogIO.from_config and register wasb remote
logging scheme (#70301)
a472c83d715 is described below
commit a472c83d715bc2c81bcfa2c19c5933e5fe7dba63
Author: Andrew Chang <[email protected]>
AuthorDate: Sat Aug 1 18:04:28 2026 +0800
Add WasbRemoteLogIO.from_config and register wasb remote logging scheme
(#70301)
Core resolves remote log handlers by URL scheme through ProvidersManager
dispatch (#67056); s3 and cloudwatch already migrated. This moves wasb onto the
same path, so Azure Blob remote logging is built by the provider's
from_config() instead of the hardcoded branch in airflow_local_settings.py.
Existing wasb:// configs resolve to an equivalent handler, and a from_config
failure falls back to the legacy path, so behaviour is unchanged.
Part of #70265. closes #70268.
---
providers/microsoft/azure/provider.yaml | 4 ++
.../providers/microsoft/azure/get_provider_info.py | 6 ++
.../microsoft/azure/log/wasb_task_handler.py | 29 ++++++++
.../microsoft/azure/log/test_wasb_task_handler.py | 83 ++++++++++++++++++++++
4 files changed, 122 insertions(+)
diff --git a/providers/microsoft/azure/provider.yaml
b/providers/microsoft/azure/provider.yaml
index d4a5e84a2ff..d1a471c3995 100644
--- a/providers/microsoft/azure/provider.yaml
+++ b/providers/microsoft/azure/provider.yaml
@@ -988,6 +988,10 @@ secrets-backends:
logging:
- airflow.providers.microsoft.azure.log.wasb_task_handler.WasbTaskHandler
+remote-logging:
+ - classpath:
airflow.providers.microsoft.azure.log.wasb_task_handler.WasbRemoteLogIO
+ scheme: wasb
+
extra-links:
-
airflow.providers.microsoft.azure.operators.data_factory.AzureDataFactoryPipelineRunLink
-
airflow.providers.microsoft.azure.operators.synapse.AzureSynapsePipelineRunLink
diff --git
a/providers/microsoft/azure/src/airflow/providers/microsoft/azure/get_provider_info.py
b/providers/microsoft/azure/src/airflow/providers/microsoft/azure/get_provider_info.py
index bf7ab5e96a8..34509081d0e 100644
---
a/providers/microsoft/azure/src/airflow/providers/microsoft/azure/get_provider_info.py
+++
b/providers/microsoft/azure/src/airflow/providers/microsoft/azure/get_provider_info.py
@@ -957,6 +957,12 @@ def get_provider_info():
],
"secrets-backends":
["airflow.providers.microsoft.azure.secrets.key_vault.AzureKeyVaultBackend"],
"logging":
["airflow.providers.microsoft.azure.log.wasb_task_handler.WasbTaskHandler"],
+ "remote-logging": [
+ {
+ "classpath":
"airflow.providers.microsoft.azure.log.wasb_task_handler.WasbRemoteLogIO",
+ "scheme": "wasb",
+ }
+ ],
"extra-links": [
"airflow.providers.microsoft.azure.operators.data_factory.AzureDataFactoryPipelineRunLink",
"airflow.providers.microsoft.azure.operators.synapse.AzureSynapsePipelineRunLink",
diff --git
a/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py
b/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py
index 5c631f4d612..51b7f7813c2 100644
---
a/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py
+++
b/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py
@@ -17,6 +17,7 @@
# under the License.
from __future__ import annotations
+import inspect
import os
import shutil
from functools import cached_property
@@ -50,6 +51,34 @@ class WasbRemoteLogIO(LoggingMixin): # noqa: D101
processors = ()
+ @classmethod
+ def from_config(cls) -> WasbRemoteLogIO:
+ """Build the remote log IO from Airflow logging configuration."""
+ remote_task_handler_kwargs = conf.getjson("logging",
"remote_task_handler_kwargs", fallback={})
+ if not isinstance(remote_task_handler_kwargs, dict):
+ raise ValueError(
+ "logging/remote_task_handler_kwargs must be a JSON object (a
python dict), we got "
+ f"{type(remote_task_handler_kwargs)}"
+ )
+ fth_params =
frozenset(inspect.signature(FileTaskHandler.__init__).parameters) - {
+ "self",
+ "base_log_folder",
+ }
+ io_kwargs = {k: v for k, v in remote_task_handler_kwargs.items() if k
not in fth_params}
+ return cls(
+ **{
+ "base_log_folder":
os.path.expanduser(conf.get_mandatory_value("logging", "base_log_folder")),
+ "remote_base": conf.get_mandatory_value("logging",
"remote_base_log_folder").removeprefix(
+ "wasb://"
+ ),
+ "delete_local_copy": conf.getboolean("logging",
"delete_local_logs"),
+ "wasb_container": conf.get_mandatory_value(
+ "azure_remote_logging", "remote_wasb_log_container",
fallback="airflow-logs"
+ ),
+ }
+ | io_kwargs,
+ )
+
def upload(self, path: str | os.PathLike, ti: RuntimeTI | None = None) ->
None:
"""Upload the given log path to the remote storage."""
path = Path(path)
diff --git
a/providers/microsoft/azure/tests/unit/microsoft/azure/log/test_wasb_task_handler.py
b/providers/microsoft/azure/tests/unit/microsoft/azure/log/test_wasb_task_handler.py
index 349ca250cbd..88f3be1c68d 100644
---
a/providers/microsoft/azure/tests/unit/microsoft/azure/log/test_wasb_task_handler.py
+++
b/providers/microsoft/azure/tests/unit/microsoft/azure/log/test_wasb_task_handler.py
@@ -41,6 +41,89 @@ pytestmark = pytest.mark.db_test
DEFAULT_DATE = datetime(2020, 8, 10)
+class TestWasbRemoteLogIOFromConfig:
+ @conf_vars(
+ {
+ ("logging", "base_log_folder"): "~/airflow/logs",
+ ("logging", "remote_base_log_folder"): "wasb://path/to/logs",
+ ("logging", "delete_local_logs"): "True",
+ ("azure_remote_logging", "remote_wasb_log_container"):
"my-container",
+ }
+ )
+ def test_from_config(self):
+ subject = WasbRemoteLogIO.from_config()
+
+ assert subject.remote_base == "path/to/logs"
+ assert subject.base_log_folder ==
Path(os.path.expanduser("~/airflow/logs"))
+ assert subject.delete_local_copy is True
+ assert subject.wasb_container == "my-container"
+
+ @conf_vars(
+ {
+ ("logging", "base_log_folder"): "/tmp/airflow/logs",
+ ("logging", "remote_base_log_folder"): "wasb://path/to/logs",
+ ("logging", "delete_local_logs"): "False",
+ ("logging", "remote_task_handler_kwargs"): '{"delete_local_copy":
true, "max_bytes": 1024}',
+ }
+ )
+ def
test_from_config_applies_io_kwargs_and_filters_file_handler_kwargs(self):
+ subject = WasbRemoteLogIO.from_config()
+
+ assert subject.delete_local_copy is True
+ assert not hasattr(subject, "max_bytes")
+ assert subject.wasb_container == "airflow-logs"
+
+ @conf_vars({("logging", "remote_task_handler_kwargs"): '["not", "a",
"dict"]'})
+ def test_from_config_rejects_non_dict_remote_task_handler_kwargs(self):
+ with pytest.raises(ValueError, match="remote_task_handler_kwargs"):
+ WasbRemoteLogIO.from_config()
+
+ def test_provider_registers_wasb_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("wasb")
+
+ assert info is not None
+ assert info.classpath ==
"airflow.providers.microsoft.azure.log.wasb_task_handler.WasbRemoteLogIO"
+
+ @pytest.mark.parametrize(
+ "manager_classpath",
+ [
+ pytest.param("airflow.providers_manager.ProvidersManager",
id="core"),
+ pytest.param(
+
"airflow.sdk.providers_manager_runtime.ProvidersManagerTaskRuntime",
id="task-runtime"
+ ),
+ ],
+ )
+ @conf_vars(
+ {
+ ("logging", "remote_logging"): "True",
+ ("logging", "remote_base_log_folder"): "wasb://path/to/logs",
+ ("logging", "remote_log_conn_id"): "wasb_default",
+ }
+ )
+ def
test_resolve_remote_task_log_uses_provider_dispatch_not_local_settings(self,
manager_classpath):
+ factory = pytest.importorskip("airflow._shared.logging.factory")
+ from airflow._shared.module_loading import import_string
+ from airflow.configuration import conf
+
+ with mock.patch.object(factory, "discover_remote_log_handler",
autospec=True) as legacy_discover:
+ remote_task_log, conn_id = factory.resolve_remote_task_log(
+ conf=conf,
+ providers_manager=import_string(manager_classpath)(),
+ import_string=import_string,
+ )
+
+ assert isinstance(remote_task_log, WasbRemoteLogIO)
+ assert remote_task_log.remote_base == "path/to/logs"
+ assert conn_id == "wasb_default"
+ legacy_discover.assert_not_called()
+
+
class TestWasbTaskHandler:
@pytest.fixture(autouse=True)
def ti(self, create_task_instance, create_log_template):