Skip to content

Commit 5747bf2

Browse files
xBis7potiukamoghrajesh
authored andcommitted
Move the traces and metrics code under a common observability package (apache#56187)
- Create shared/observability with base metrics and traces implementations - Add task-sdk observability module with own Stats and Trace wrappers - task SDK no longer imports from airflow-core observability - Each component reads its own configuration - Export only Trace in public API (Stats is internal) --------- Co-authored-by: Jarek Potiuk <jarek@potiuk.com> Co-authored-by: Amogh Desai <amoghrajesh1999@gmail.com>
1 parent 367d9a3 commit 5747bf2

103 files changed

Lines changed: 1440 additions & 502 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

airflow-core/pyproject.toml

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -227,6 +227,7 @@ exclude = [
227227
"../shared/secrets_masker/src/airflow_shared/secrets_masker" = "src/airflow/_shared/secrets_masker"
228228
"../shared/timezones/src/airflow_shared/timezones" = "src/airflow/_shared/timezones"
229229
"../shared/secrets_backend/src/airflow_shared/secrets_backend" = "src/airflow/_shared/secrets_backend"
230+
"../shared/observability/src/airflow_shared/observability" = "src/airflow/_shared/observability"
230231

231232
[tool.hatch.build.targets.custom]
232233
path = "./hatch_build.py"
@@ -299,4 +300,5 @@ shared_distributions = [
299300
"apache-airflow-shared-secrets-backend",
300301
"apache-airflow-shared-secrets-masker",
301302
"apache-airflow-shared-timezones",
303+
"apache-airflow-shared-observability",
302304
]

airflow-core/src/airflow/__init__.py

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -89,6 +89,10 @@
8989
# Deprecated lazy imports
9090
"AirflowException": (".exceptions", "AirflowException", True),
9191
"Dataset": (".sdk", "Asset", True),
92+
"Stats": (".observability.stats", "Stats", True),
93+
"Trace": (".observability.trace", "Trace", True),
94+
"metrics": (".observability.metrics", "", True),
95+
"traces": (".observability.traces", "", True),
9296
}
9397
if TYPE_CHECKING:
9498
# These objects are imported by PEP-562, however, static analyzers and IDE's
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
../../../../shared/observability/src/airflow_shared/observability

airflow-core/src/airflow/assets/manager.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,7 @@
3838
DagScheduleAssetUriReference,
3939
PartitionedAssetKeyLog,
4040
)
41-
from airflow.stats import Stats
41+
from airflow.observability.stats import Stats
4242
from airflow.utils.log.logging_mixin import LoggingMixin
4343
from airflow.utils.sqlalchemy import get_dialect_name
4444

airflow-core/src/airflow/cli/commands/daemon_utils.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -75,7 +75,7 @@ def run_command_with_daemon_option(
7575

7676
with ctx:
7777
# in daemon context stats client needs to be reinitialized.
78-
from airflow.stats import Stats
78+
from airflow.observability.stats import Stats
7979

8080
Stats.instance = None
8181
callback()

airflow-core/src/airflow/dag_processing/manager.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -58,10 +58,10 @@
5858
from airflow.models.dagwarning import DagWarning
5959
from airflow.models.db_callback_request import DbCallbackRequest
6060
from airflow.models.errors import ParseImportError
61+
from airflow.observability.stats import Stats
62+
from airflow.observability.trace import DebugTrace
6163
from airflow.sdk import SecretCache
6264
from airflow.sdk.log import init_log_file, logging_processors
63-
from airflow.stats import Stats
64-
from airflow.traces.tracer import DebugTrace
6565
from airflow.utils.file import list_py_file_paths, might_contain_dag
6666
from airflow.utils.log.logging_mixin import LoggingMixin
6767
from airflow.utils.net import get_hostname

airflow-core/src/airflow/dag_processing/processor.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,7 @@
3636
)
3737
from airflow.configuration import conf
3838
from airflow.dag_processing.dagbag import DagBag
39+
from airflow.observability.stats import Stats
3940
from airflow.sdk.exceptions import TaskNotFound
4041
from airflow.sdk.execution_time.comms import (
4142
ConnectionResult,
@@ -66,7 +67,6 @@
6667
from airflow.sdk.execution_time.supervisor import WatchedSubprocess
6768
from airflow.sdk.execution_time.task_runner import RuntimeTaskInstance, _send_error_email_notification
6869
from airflow.serialization.serialized_objects import LazyDeserializedDAG, SerializedDAG
69-
from airflow.stats import Stats
7070
from airflow.utils.file import iter_airflow_imports
7171
from airflow.utils.state import TaskInstanceState
7272

airflow-core/src/airflow/executors/base_executor.py

Lines changed: 5 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -27,15 +27,14 @@
2727

2828
import pendulum
2929

30+
from airflow._shared.observability.traces import NO_TRACE_ID
3031
from airflow.cli.cli_config import DefaultHelpParser
3132
from airflow.configuration import conf
3233
from airflow.executors import workloads
3334
from airflow.executors.executor_loader import ExecutorLoader
3435
from airflow.models import Log
35-
from airflow.stats import Stats
36-
from airflow.traces import NO_TRACE_ID
37-
from airflow.traces.tracer import DebugTrace, Trace, add_debug_span, gen_context
38-
from airflow.traces.utils import gen_span_id_from_ti_key
36+
from airflow.observability.stats import Stats
37+
from airflow.observability.trace import DebugTrace, Trace, add_debug_span
3938
from airflow.utils.log.logging_mixin import LoggingMixin
4039
from airflow.utils.state import TaskInstanceState
4140
from airflow.utils.thread_safe_dict import ThreadSafeDict
@@ -419,11 +418,9 @@ def fail(self, key: TaskInstanceKey, info=None) -> None:
419418
"""
420419
trace_id = Trace.get_current_span().get_span_context().trace_id
421420
if trace_id != NO_TRACE_ID:
422-
span_id = int(gen_span_id_from_ti_key(key, as_int=True))
423-
with DebugTrace.start_span(
421+
with DebugTrace.start_child_span(
424422
span_name="fail",
425423
component="BaseExecutor",
426-
parent_sc=gen_context(trace_id=trace_id, span_id=span_id),
427424
) as span:
428425
span.set_attributes(
429426
{
@@ -446,11 +443,9 @@ def success(self, key: TaskInstanceKey, info=None) -> None:
446443
"""
447444
trace_id = Trace.get_current_span().get_span_context().trace_id
448445
if trace_id != NO_TRACE_ID:
449-
span_id = int(gen_span_id_from_ti_key(key, as_int=True))
450-
with DebugTrace.start_span(
446+
with DebugTrace.start_child_span(
451447
span_name="success",
452448
component="BaseExecutor",
453-
parent_sc=gen_context(trace_id=trace_id, span_id=span_id),
454449
) as span:
455450
span.set_attributes(
456451
{

airflow-core/src/airflow/jobs/dag_processor_job_runner.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@
2121

2222
from airflow.jobs.base_job_runner import BaseJobRunner
2323
from airflow.jobs.job import Job, perform_heartbeat
24-
from airflow.stats import Stats
24+
from airflow.observability.stats import Stats
2525
from airflow.utils.log.logging_mixin import LoggingMixin
2626

2727
if TYPE_CHECKING:

airflow-core/src/airflow/jobs/job.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -35,8 +35,8 @@
3535
from airflow.executors.executor_loader import ExecutorLoader
3636
from airflow.listeners.listener import get_listener_manager
3737
from airflow.models.base import ID_LEN, Base
38-
from airflow.stats import Stats
39-
from airflow.traces.tracer import DebugTrace, add_debug_span
38+
from airflow.observability.stats import Stats
39+
from airflow.observability.trace import DebugTrace, add_debug_span
4040
from airflow.utils.helpers import convert_camel_to_snake
4141
from airflow.utils.log.logging_mixin import LoggingMixin
4242
from airflow.utils.net import get_hostname

0 commit comments

Comments
 (0)