Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion airflow/cli/commands/triggerer_command.py
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,8 @@ def triggerer(args):
"""Starts Airflow Triggerer."""
settings.MASK_SECRETS_IN_LOGS = True
print(settings.HEADER)
triggerer_job_runner = TriggererJobRunner(job=Job(), capacity=args.capacity)
triggerer_heartrate = conf.getfloat("triggerer", "JOB_HEARTBEAT_SEC")
triggerer_job_runner = TriggererJobRunner(job=Job(heartrate=triggerer_heartrate), capacity=args.capacity)

if args.daemon:
pid, stdout, stderr, log_file = setup_locations(
Expand Down
7 changes: 7 additions & 0 deletions airflow/config_templates/config.yml
Original file line number Diff line number Diff line change
Expand Up @@ -2558,6 +2558,13 @@ triggerer:
type: string
example: ~
default: "1000"
job_heartbeat_sec:
description: |
How often to heartbeat the Triggerer job to ensure it hasn't been killed.
version_added: 2.6.3

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Since we are adding a new configuration, this should be moved to 2.7.0

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not necessarily. I think we do not have a rule "do not add new configuration options" in patchlevel. We have the rule of not adding new features, but new configuration might be added in order to implement a bugfix (which I think is the case here).

I think it's OK (and we've done that in the past) that we introduced new configuration in patchlevel version if they served bugfix purpose. There are a few configs already that were added in non- .0 versions (smtp ones for example).

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

(and yes, that was also my initial reaction, but I thought about it and I looked at the past configuration entries we created and it's not obvious that "new configuration == new feature".

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ephraimbuddy Do you agree with @potiuk, or should I move it to 2.7.0?

type: float
example: ~
default: "5"
kerberos:
description: ~
options:
Expand Down
3 changes: 3 additions & 0 deletions airflow/config_templates/default_airflow.cfg
Original file line number Diff line number Diff line change
Expand Up @@ -1326,6 +1326,9 @@ task_queued_timeout_check_interval = 120.0
# How many triggers a single Triggerer will run at once, by default.
default_capacity = 1000

# How often to heartbeat the Triggerer job to ensure it hasn't been killed.
job_heartbeat_sec = 5

[kerberos]
ccache = /tmp/airflow_krb5_ccache

Expand Down
2 changes: 1 addition & 1 deletion airflow/jobs/triggerer_job_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -371,7 +371,7 @@ def load_triggers(self):
adds them to our runner, and then removes ones from it we no longer
need.
"""
Trigger.assign_unassigned(self.job.id, self.capacity)
Trigger.assign_unassigned(self.job.id, self.capacity, self.job.heartrate)
ids = Trigger.ids_for_triggerer(self.job.id)
self.trigger_runner.update_triggers(set(ids))

Expand Down
13 changes: 8 additions & 5 deletions airflow/models/trigger.py
Original file line number Diff line number Diff line change
Expand Up @@ -200,10 +200,11 @@ def ids_for_triggerer(cls, triggerer_id, session: Session = NEW_SESSION) -> list
@classmethod
@internal_api_call
@provide_session
def assign_unassigned(cls, triggerer_id, capacity, session: Session = NEW_SESSION) -> None:
def assign_unassigned(cls, triggerer_id, capacity, heartrate, session: Session = NEW_SESSION) -> None:
"""
Takes a triggerer_id and the capacity for that triggerer and assigns unassigned
triggers until that capacity is reached, or there are no more unassigned triggers.
Takes a triggerer_id, the capacity for that triggerer and the Triggerer job heartrate,
and assigns unassigned triggers until that capacity is reached, or there are no more
unassigned triggers.
"""
from airflow.jobs.job import Job # To avoid circular import

Expand All @@ -212,12 +213,14 @@ def assign_unassigned(cls, triggerer_id, capacity, session: Session = NEW_SESSIO

if capacity <= 0:
return

# we multiply heartrate by a grace_multiplier to give the triggerer
# a chance to heartbeat before we consider it dead
health_check_threshold = heartrate * 2.1
Comment thread
hussein-awala marked this conversation as resolved.
alive_triggerer_ids = [
row[0]
for row in session.query(Job.id).filter(
Job.end_date.is_(None),
Job.latest_heartbeat > timezone.utcnow() - datetime.timedelta(seconds=30),
Job.latest_heartbeat > timezone.utcnow() - datetime.timedelta(seconds=health_check_threshold),
Job.job_type == "TriggererJob",
)
]
Expand Down
64 changes: 60 additions & 4 deletions tests/models/test_trigger.py
Original file line number Diff line number Diff line change
Expand Up @@ -141,16 +141,17 @@ def test_assign_unassigned(session, create_task_instance):
"""
Tests that unassigned triggers of all appropriate states are assigned.
"""
finished_triggerer = Job(heartrate=10, state=State.SUCCESS)
triggerer_heartrate = 10
finished_triggerer = Job(heartrate=triggerer_heartrate, state=State.SUCCESS)
TriggererJobRunner(finished_triggerer)
finished_triggerer.end_date = timezone.utcnow() - datetime.timedelta(hours=1)
session.add(finished_triggerer)
assert not finished_triggerer.is_alive()
healthy_triggerer = Job(heartrate=10, state=State.RUNNING)
healthy_triggerer = Job(heartrate=triggerer_heartrate, state=State.RUNNING)
TriggererJobRunner(healthy_triggerer)
session.add(healthy_triggerer)
assert healthy_triggerer.is_alive()
new_triggerer = Job(heartrate=10, state=State.RUNNING)
new_triggerer = Job(heartrate=triggerer_heartrate, state=State.RUNNING)
TriggererJobRunner(new_triggerer)
session.add(new_triggerer)
assert new_triggerer.is_alive()
Expand All @@ -169,7 +170,7 @@ def test_assign_unassigned(session, create_task_instance):
session.add(trigger_unassigned_to_triggerer)
session.commit()
assert session.query(Trigger).count() == 3
Trigger.assign_unassigned(new_triggerer.id, 100, session=session)
Trigger.assign_unassigned(new_triggerer.id, 100, session=session, heartrate=triggerer_heartrate)
session.expire_all()
# Check that trigger on killed triggerer and unassigned trigger are assigned to new triggerer
assert (
Expand All @@ -187,6 +188,61 @@ def test_assign_unassigned(session, create_task_instance):
)


@pytest.mark.parametrize("check_triggerer_heartrate", [10, 60, 300])
def test_assign_unassigned_missing_heartbeat(session, create_task_instance, check_triggerer_heartrate):
"""
Tests that the triggers assigned to a dead triggers are considered as unassigned
and they are assigned to an alive triggerer.
"""
import time_machine

block_triggerer_heartrate = 9999
with time_machine.travel(datetime.datetime.utcnow(), tick=False) as t:
first_triggerer = Job(heartrate=block_triggerer_heartrate, state=State.RUNNING)
TriggererJobRunner(first_triggerer)
session.add(first_triggerer)
assert first_triggerer.is_alive()
second_triggerer = Job(heartrate=block_triggerer_heartrate, state=State.RUNNING)
TriggererJobRunner(second_triggerer)
session.add(second_triggerer)
assert second_triggerer.is_alive()
session.commit()
trigger_on_first_triggerer = Trigger(classpath="airflow.triggers.testing.SuccessTrigger", kwargs={})
trigger_on_first_triggerer.id = 1
trigger_on_first_triggerer.triggerer_id = first_triggerer.id
trigger_on_second_triggerer = Trigger(classpath="airflow.triggers.testing.SuccessTrigger", kwargs={})
trigger_on_second_triggerer.id = 2
trigger_on_second_triggerer.triggerer_id = second_triggerer.id
session.add(trigger_on_first_triggerer)
session.add(trigger_on_second_triggerer)
session.commit()
assert session.query(Trigger).count() == 2
triggers_ids = [
(first_triggerer.id, second_triggerer.id),
(first_triggerer.id, second_triggerer.id),
(first_triggerer.id, second_triggerer.id),
# Check that after more than 2.1 heartrates, the first triggerer is considered dead
# and the first trigger is assigned to the second triggerer
(second_triggerer.id, second_triggerer.id),
]
for i in range(4):
Trigger.assign_unassigned(
second_triggerer.id, 100, session=session, heartrate=check_triggerer_heartrate
)
session.expire_all()
# Check that trigger on killed triggerer and unassigned trigger are assigned to new triggerer
assert (
session.query(Trigger).filter(Trigger.id == trigger_on_first_triggerer.id).one().triggerer_id
== triggers_ids[i][0]
)
assert (
session.query(Trigger).filter(Trigger.id == trigger_on_second_triggerer.id).one().triggerer_id
== triggers_ids[i][1]
)
t.shift(datetime.timedelta(seconds=check_triggerer_heartrate))
second_triggerer.latest_heartbeat += datetime.timedelta(seconds=check_triggerer_heartrate)


def test_get_sorted_triggers(session, create_task_instance):
"""
Tests that triggers are sorted by the creation_date.
Expand Down