kaxil commented on code in PR #74204:
URL: https://github.com/apache/airflow/pull/74204#discussion_r4183703253
##########
providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_task_handler.py:
##########
@@ -554,6 +581,7 @@ def set_context(self, ti: TaskInstance, *, identifier: str
| None = None) -> Non
),
"try_number": str(ti.try_number),
"log_id": _render_log_id(self.log_id_template, ti,
ti.try_number),
+ **_get_ti_id_fields(ti),
Review Comment:
Can this ever add anything? On Airflow 3 nothing calls handler `set_context`
any more (the only caller is the module-level `set_context` in
`logging_mixin.py`, which has no callers in core, the SDK, or providers), and
on Airflow 2 the TI has no `id` so this returns `{}`. Same for
`os_task_handler.py:496`, where the provider floor already excludes Airflow 2.
I'd drop both.
##########
providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_task_handler.py:
##########
@@ -226,6 +230,29 @@ def _render_log_id(log_id_template: str, ti: TaskInstance
| TaskInstanceKey, try
)
+def _get_ti_id_fields(ti: TaskInstance | TaskInstanceKey) -> dict[str, str]:
+ # Airflow 2 task instances have no id.
+ return {"ti_id": str(ti_id)} if (ti_id := getattr(ti, "id", None)) else {}
+
+
+def _build_log_query(log_id: str, ti: RuntimeTI) -> list[dict[str, Any]]:
+ log_id_match = {"match_phrase": {"log_id": log_id}}
+ # Before 3.4 a cleared task instance gets a new id, which can differ from
the id its logs were written under.
+ if not AIRFLOW_V_3_4_PLUS:
+ return [log_id_match]
+ return [
+ {
+ "bool": {
+ "should": [
+ {"match_phrase": {"ti_id": str(ti.id)}},
Review Comment:
On 3.4 this matches on the id alone, but 3.3.x could reuse one id across two
tries. In 3.3.2 `handle_failure` only calls `prepare_db_for_next_try` when the
TI is RUNNING, so a deferrable or reschedule task that uploaded part of try N
under id A, then failed in QUEUED after resuming (pod never started, say),
retries as try N+1 still under A. After upgrading to 3.4 the try N+1 view pulls
in try N's documents too, and since both start at offset 1 the lines
interleave. Would ANDing `log_id_match` with the `should` (and the same with
`try_number` for the GCP filter) work for you? Loop passes sharing coordinates
would still be split by `ti_id`. Unless matching on id alone is deliberate,
e.g. to survive a `log_id_template` change. Same in `os_task_handler.py:285`
and `cloud_logging_task_handler.py:281`.
##########
providers/google/src/airflow/providers/google/cloud/log/cloud_logging_task_handler.py:
##########
@@ -269,8 +272,18 @@ def escape_label_value(value: str) -> str:
for key, value in self.resource.labels.items():
log_filters.append(f"resource.labels.{escape_label_key(key)}={escape_label_value(value)}")
- for key, value in ti_labels.items():
-
log_filters.append(f"labels.{escape_label_key(key)}={escape_label_value(value)}")
+ label_conditions = {
+ key: f"labels.{escape_label_key(key)}={escape_label_value(value)}"
+ for key, value in ti_labels.items()
+ }
+ ti_id_condition = label_conditions.pop(LABEL_TI_ID, None)
+ # Before 3.4 a cleared task instance gets a new id, which can differ
from the id its logs were written under.
+ if ti_id_condition and AIRFLOW_V_3_4_PLUS:
+ # Entries written before the ``ti_id`` label existed are matched
by the other labels.
Review Comment:
FYI, pre-existing and not from this PR: I don't think this holds for Airflow
3 entries. The structlog writer path above never sets a `logical_date` label,
but `_task_instance_to_labels` always puts it in the filter, so the legacy
branch can't match anything that writer produced. The new `ti_id` branch is
what makes them readable after this change.
##########
providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_task_handler.py:
##########
@@ -801,6 +832,7 @@ def _parse_raw_log(self, log: str, log_id: str) ->
list[dict[str, Any]]:
log_dict.update(
{
"log_id": log_id,
+ **(extra_fields or {}),
Review Comment:
Worth adding `ti_id` to the optional fields in the document schema section
of `docs/logging/index.rst`? Once the id branch is active, whether an
externally shipped document carries it changes which branch of the query it
matches, and Fluent Bit / Logstash setups will forward it from the task's JSON
lines without anyone deciding to.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]