Skip to content

Commit 60d4bcd

Browse files
authored
Revert "Enable individual trigger logging (#27758)" (#29472)
This reverts commit 1b18a50.
1 parent c44b24e commit 60d4bcd

54 files changed

Lines changed: 487 additions & 2383 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/api_connexion/endpoints/log_endpoint.py

Lines changed: 2 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,6 @@
2121
from flask import Response, request
2222
from itsdangerous.exc import BadSignature
2323
from itsdangerous.url_safe import URLSafeSerializer
24-
from sqlalchemy.orm import joinedload
2524
from sqlalchemy.orm.session import Session
2625

2726
from airflow.api_connexion import security
@@ -74,10 +73,9 @@ def get_log(
7473
metadata["download_logs"] = False
7574

7675
task_log_reader = TaskLogReader()
77-
7876
if not task_log_reader.supports_read:
7977
raise BadRequest("Task log handler does not support read logs.")
80-
query = (
78+
ti = (
8179
session.query(TaskInstance)
8280
.filter(
8381
TaskInstance.task_id == task_id,
@@ -86,10 +84,8 @@ def get_log(
8684
TaskInstance.map_index == map_index,
8785
)
8886
.join(TaskInstance.dag_run)
89-
.options(joinedload("trigger"))
90-
.options(joinedload("trigger.triggerer_job"))
87+
.one_or_none()
9188
)
92-
ti = query.one_or_none()
9389
if ti is None:
9490
metadata["end_of_log"] = True
9591
raise NotFound(title="TaskInstance not found")

airflow/cli/cli_parser.py

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2115,7 +2115,6 @@ class GroupCommand(NamedTuple):
21152115
ARG_LOG_FILE,
21162116
ARG_CAPACITY,
21172117
ARG_VERBOSE,
2118-
ARG_SKIP_SERVE_LOGS,
21192118
),
21202119
),
21212120
ActionCommand(

airflow/cli/commands/task_command.py

Lines changed: 11 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -44,7 +44,6 @@
4444
from airflow.models.dag import DAG
4545
from airflow.models.dagrun import DagRun
4646
from airflow.models.operator import needs_expansion
47-
from airflow.models.taskinstance import TaskReturnCode
4847
from airflow.settings import IS_K8S_EXECUTOR_POD
4948
from airflow.ti_deps.dep_context import DepContext
5049
from airflow.ti_deps.dependencies_deps import SCHEDULER_QUEUED_DEPS
@@ -59,7 +58,6 @@
5958
suppress_logs_and_warning,
6059
)
6160
from airflow.utils.dates import timezone
62-
from airflow.utils.log.file_task_handler import _set_task_deferred_context_var
6361
from airflow.utils.log.logging_mixin import StreamLogWriter
6462
from airflow.utils.log.secrets_masker import RedactedIO
6563
from airflow.utils.net import get_hostname
@@ -184,7 +182,7 @@ def _get_ti(
184182
return ti, dr_created
185183

186184

187-
def _run_task_by_selected_method(args, dag: DAG, ti: TaskInstance) -> None | TaskReturnCode:
185+
def _run_task_by_selected_method(args, dag: DAG, ti: TaskInstance) -> None:
188186
"""
189187
Runs the task based on a mode.
190188
@@ -195,11 +193,11 @@ def _run_task_by_selected_method(args, dag: DAG, ti: TaskInstance) -> None | Tas
195193
- by executor
196194
"""
197195
if args.local:
198-
return _run_task_by_local_task_job(args, ti)
196+
_run_task_by_local_task_job(args, ti)
199197
elif args.raw:
200-
return _run_raw_task(args, ti)
198+
_run_raw_task(args, ti)
201199
else:
202-
return _run_task_by_executor(args, dag, ti)
200+
_run_task_by_executor(args, dag, ti)
203201

204202

205203
def _run_task_by_executor(args, dag, ti):
@@ -241,7 +239,7 @@ def _run_task_by_executor(args, dag, ti):
241239
executor.end()
242240

243241

244-
def _run_task_by_local_task_job(args, ti) -> TaskReturnCode | None:
242+
def _run_task_by_local_task_job(args, ti):
245243
"""Run LocalTaskJob, which monitors the raw task execution process."""
246244
run_job = LocalTaskJob(
247245
task_instance=ti,
@@ -256,14 +254,11 @@ def _run_task_by_local_task_job(args, ti) -> TaskReturnCode | None:
256254
external_executor_id=_extract_external_executor_id(args),
257255
)
258256
try:
259-
ret = run_job.run()
257+
run_job.run()
260258

261259
finally:
262260
if args.shut_down_logging:
263261
logging.shutdown()
264-
with suppress(ValueError):
265-
return TaskReturnCode(ret)
266-
return None
267262

268263

269264
RAW_TASK_UNSUPPORTED_OPTION = [
@@ -274,9 +269,9 @@ def _run_task_by_local_task_job(args, ti) -> TaskReturnCode | None:
274269
]
275270

276271

277-
def _run_raw_task(args, ti: TaskInstance) -> None | TaskReturnCode:
272+
def _run_raw_task(args, ti: TaskInstance) -> None:
278273
"""Runs the main task handling code."""
279-
return ti._run_raw_task(
274+
ti._run_raw_task(
280275
mark_success=args.mark_success,
281276
job_id=args.job_id,
282277
pool=args.pool,
@@ -412,21 +407,18 @@ def task_run(args, dag=None):
412407
# this should be last thing before running, to reduce likelihood of an open session
413408
# which can cause trouble if running process in a fork.
414409
settings.reconfigure_orm(disable_connection_pool=True)
415-
task_return_code = None
410+
416411
try:
417412
if args.interactive:
418-
task_return_code = _run_task_by_selected_method(args, dag, ti)
413+
_run_task_by_selected_method(args, dag, ti)
419414
else:
420415
with _move_task_handlers_to_root(ti), _redirect_stdout_to_ti_log(ti):
421-
task_return_code = _run_task_by_selected_method(args, dag, ti)
422-
if task_return_code == TaskReturnCode.DEFERRED:
423-
_set_task_deferred_context_var()
416+
_run_task_by_selected_method(args, dag, ti)
424417
finally:
425418
try:
426419
get_listener_manager().hook.before_stopping(component=TaskCommandMarker())
427420
except Exception:
428421
pass
429-
return task_return_code
430422

431423

432424
@cli_utils.action_cli(check_db=False)

airflow/cli/commands/triggerer_command.py

Lines changed: 4 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -18,35 +18,14 @@
1818
from __future__ import annotations
1919

2020
import signal
21-
from contextlib import contextmanager
22-
from functools import partial
23-
from multiprocessing import Process
24-
from typing import Generator
2521

2622
import daemon
2723
from daemon.pidfile import TimeoutPIDLockFile
2824

2925
from airflow import settings
30-
from airflow.configuration import conf
3126
from airflow.jobs.triggerer_job import TriggererJob
3227
from airflow.utils import cli as cli_utils
3328
from airflow.utils.cli import setup_locations, setup_logging, sigint_handler, sigquit_handler
34-
from airflow.utils.serve_logs import serve_logs
35-
36-
37-
@contextmanager
38-
def _serve_logs(skip_serve_logs: bool = False) -> Generator[None, None, None]:
39-
"""Starts serve_logs sub-process"""
40-
sub_proc = None
41-
if skip_serve_logs is False:
42-
port = conf.getint("logging", "trigger_log_server_port", fallback=8794)
43-
sub_proc = Process(target=partial(serve_logs, port=port))
44-
sub_proc.start()
45-
try:
46-
yield
47-
finally:
48-
if sub_proc:
49-
sub_proc.terminate()
5029

5130

5231
@cli_utils.action_cli
@@ -65,18 +44,18 @@ def triggerer(args):
6544
stdout_handle.truncate(0)
6645
stderr_handle.truncate(0)
6746

68-
daemon_context = daemon.DaemonContext(
47+
ctx = daemon.DaemonContext(
6948
pidfile=TimeoutPIDLockFile(pid, -1),
7049
files_preserve=[handle],
7150
stdout=stdout_handle,
7251
stderr=stderr_handle,
7352
umask=int(settings.DAEMON_UMASK, 8),
7453
)
75-
with daemon_context, _serve_logs(args.skip_serve_logs):
54+
with ctx:
7655
job.run()
56+
7757
else:
7858
signal.signal(signal.SIGINT, sigint_handler)
7959
signal.signal(signal.SIGTERM, sigint_handler)
8060
signal.signal(signal.SIGQUIT, sigquit_handler)
81-
with _serve_logs(args.skip_serve_logs):
82-
job.run()
61+
job.run()

airflow/config_templates/config.yml

Lines changed: 0 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -788,24 +788,6 @@ logging:
788788
type: string
789789
example: ~
790790
default: "8793"
791-
trigger_log_server_port:
792-
description: |
793-
Port to serve logs from for triggerer. See worker_log_server_port description
794-
for more info.
795-
version_added: 2.6.0
796-
type: string
797-
example: ~
798-
default: "8794"
799-
interleave_timestamp_parser:
800-
description: |
801-
We must parse timestamps to interleave logs between trigger and task. To do so,
802-
we need to parse timestamps in log files. In case your log format is non-standard,
803-
you may provide import path to callable which takes a string log line and returns
804-
the timestamp (datetime.datetime compatible).
805-
version_added: 2.6.0
806-
type: string
807-
example: path.to.my_func
808-
default: ~
809791
metrics:
810792
description: |
811793
StatsD (https://github.com/etsy/statsd) integration settings.

airflow/config_templates/default_airflow.cfg

Lines changed: 0 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -434,17 +434,6 @@ extra_logger_names =
434434
# visible from the main web server to connect into the workers.
435435
worker_log_server_port = 8793
436436

437-
# Port to serve logs from for triggerer. See worker_log_server_port description
438-
# for more info.
439-
trigger_log_server_port = 8794
440-
441-
# We must parse timestamps to interleave logs between trigger and task. To do so,
442-
# we need to parse timestamps in log files. In case your log format is non-standard,
443-
# you may provide import path to callable which takes a string log line and returns
444-
# the timestamp (datetime.datetime compatible).
445-
# Example: interleave_timestamp_parser = path.to.my_func
446-
# interleave_timestamp_parser =
447-
448437
[metrics]
449438

450439
# StatsD (https://github.com/etsy/statsd) integration settings.

airflow/example_dags/example_time_delta_sensor_async.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,6 @@
3636
catchup=False,
3737
tags=["example"],
3838
) as dag:
39-
wait = TimeDeltaSensorAsync(task_id="wait", delta=datetime.timedelta(seconds=30))
39+
wait = TimeDeltaSensorAsync(task_id="wait", delta=datetime.timedelta(seconds=10))
4040
finish = EmptyOperator(task_id="finish")
4141
wait >> finish

airflow/executors/base_executor.py

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -356,15 +356,14 @@ def execute_async(
356356
"""
357357
raise NotImplementedError()
358358

359-
def get_task_log(self, ti: TaskInstance) -> tuple[list[str], list[str]]:
359+
def get_task_log(self, ti: TaskInstance, log: str = "") -> None | str | tuple[str, dict[str, bool]]:
360360
"""
361361
This method can be implemented by any child class to return the task logs.
362362
363363
:param ti: A TaskInstance object
364364
:param log: log str
365365
:return: logs or tuple of logs and meta dict
366366
"""
367-
return [], []
368367

369368
def end(self) -> None: # pragma: no cover
370369
"""Wait synchronously for the previously submitted job to complete."""

airflow/executors/celery_kubernetes_executor.py

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -140,11 +140,11 @@ def queue_task_instance(
140140
cfg_path=cfg_path,
141141
)
142142

143-
def get_task_log(self, ti: TaskInstance) -> tuple[list[str], list[str]]:
143+
def get_task_log(self, ti: TaskInstance, log: str = "") -> None | str | tuple[str, dict[str, bool]]:
144144
"""Fetch task log from Kubernetes executor"""
145145
if ti.queue == self.kubernetes_executor.kubernetes_queue:
146-
return self.kubernetes_executor.get_task_log(ti=ti)
147-
return [], []
146+
return self.kubernetes_executor.get_task_log(ti=ti, log=log)
147+
return None
148148

149149
def has_task(self, task_instance: TaskInstance) -> bool:
150150
"""

airflow/executors/kubernetes_executor.py

Lines changed: 10 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -781,16 +781,14 @@ def _get_pod_namespace(ti: TaskInstance):
781781
namespace = pod_override.metadata.namespace
782782
return namespace or conf.get("kubernetes_executor", "namespace", fallback="default")
783783

784-
def get_task_log(self, ti: TaskInstance) -> tuple[list[str], list[str]]:
785-
messages = []
786-
log = []
784+
def get_task_log(self, ti: TaskInstance, log: str = "") -> str | tuple[str, dict[str, bool]]:
785+
787786
try:
788-
from airflow.kubernetes.kube_client import get_kube_client
789787
from airflow.kubernetes.pod_generator import PodGenerator
790788

791789
client = get_kube_client()
792790

793-
messages.append(f"Trying to get logs (last 100 lines) from worker pod {ti.hostname}")
791+
log += f"*** Trying to get logs (last 100 lines) from worker pod {ti.hostname} ***\n\n"
794792
selector = PodGenerator.build_selector_for_k8s_executor_pod(
795793
dag_id=ti.dag_id,
796794
task_id=ti.task_id,
@@ -818,10 +816,13 @@ def get_task_log(self, ti: TaskInstance) -> tuple[list[str], list[str]]:
818816
)
819817

820818
for line in res:
821-
log.append(line.decode())
822-
except Exception as e:
823-
messages.append(f"Reading from k8s pod logs failed: {str(e)}")
824-
return messages, ["\n".join(log)]
819+
log += line.decode()
820+
821+
return log
822+
823+
except Exception as f:
824+
log += f"*** Unable to fetch logs from worker pod {ti.hostname} ***\n{str(f)}\n\n"
825+
return log, {"end_of_log": True}
825826

826827
def try_adopt_task_instances(self, tis: Sequence[TaskInstance]) -> Sequence[TaskInstance]:
827828
tis_to_flush = [ti for ti in tis if not ti.queued_by_job_id]

0 commit comments

Comments
 (0)