Skip to content

Commit eccd4a2

Browse files
authored
[Airflow integration - logs pipeline asset] Add support for new timestamp format (DataDog#22730)
* date parsing rules to supprot old and new formats * add test with timestamp in new format * keep samples as before * test all formats
1 parent a9fa87b commit eccd4a2

2 files changed

Lines changed: 59 additions & 2 deletions

File tree

airflow/assets/logs/airflow.yaml

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -38,11 +38,15 @@ pipeline:
3838
Exception: Hello world fail!
3939
grok:
4040
supportRules: |
41-
_date %{date("yyyy-MM-dd HH:mm:ss,SSS"):timestamp}
41+
_date_legacy %{date("yyyy-MM-dd HH:mm:ss,SSS"):timestamp}
42+
_date_iso_offset %{date("yyyy-MM-dd'T'HH:mm:ss.SSSZ"):timestamp}
43+
_date_iso_zulu %{date("yyyy-MM-dd'T'HH:mm:ss.SSS'Z'"):timestamp}
44+
_date (%{_date_iso_offset}|%{_date_iso_zulu}|%{_date_legacy})
45+
_date_any (%{date("yyyy-MM-dd'T'HH:mm:ss.SSSZ")}|%{date("yyyy-MM-dd'T'HH:mm:ss.SSS'Z'")}|%{date("yyyy-MM-dd HH:mm:ss,SSS")})
4246
_file %{notSpace:filename}:%{integer:lineno}
4347
_base_prefix \[%{_date}\]\s+(\{\{%{_file}\}\}|\{%{_file}\})\s+%{word:level}\s+-
4448
# Some logs have duplicate prefixes, we only collect job_id and subtask name from first prefix if present.
45-
_prefix (\[%{date("yyyy-MM-dd HH:mm:ss,SSS")}\]\s+\{\{%{notSpace}:%{integer}\}\}\s+%{word}\s+-\s+)?(Job %{integer:jobid}: Subtask %{word:subtask}\s+)?%{_base_prefix}
49+
_prefix (\[%{_date_any}\]\s+(\{\{%{notSpace}:%{integer}\}\}|\{%{notSpace}:%{integer}\})\s+%{word}\s+-\s+)?(Job %{integer:jobid}: Subtask %{word:subtask}\s+)?%{_base_prefix}
4650
4751
matchRules: |
4852
airflow_stacktrace %{_prefix}\s+%{data:message}\nTraceback \(most recent call last\):\n%{data:error.stack}

airflow/assets/logs/airflow_tests.yaml

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -214,6 +214,59 @@ tests:
214214
tags:
215215
- "source:LOGS_SOURCE"
216216
timestamp: 1577381689039
217+
-
218+
sample: "[2025-04-10T17:15:01.220Z] {logging_mixin.py:188} INFO - Running task"
219+
result:
220+
custom:
221+
filename: "logging_mixin.py"
222+
level: "INFO"
223+
lineno: 188
224+
timestamp: 1744305301220
225+
message: "Running task"
226+
status: "info"
227+
tags:
228+
- "source:LOGS_SOURCE"
229+
timestamp: 1744305301220
230+
-
231+
sample: "[2025-04-10T17:15:01.220+0100] {base_task_runner.py:115} INFO - Job 12: Subtask my_task [2025-04-10T17:15:01.450+0100] {settings.py:211} INFO - Setting up DB connection pool (PID 42)"
232+
result:
233+
custom:
234+
filename: "settings.py"
235+
jobid: 12
236+
level: "INFO"
237+
lineno: 211
238+
subtask: "my_task"
239+
timestamp: 1744301701450
240+
message: "Setting up DB connection pool (PID 42)"
241+
status: "info"
242+
tags:
243+
- "source:LOGS_SOURCE"
244+
timestamp: 1744301701450
245+
-
246+
sample: |
247+
[2025-04-10T17:15:01.220+0000] {taskinstance.py:1058} ERROR - Task failed!
248+
Traceback (most recent call last):
249+
File "/usr/local/lib/python3.11/site-packages/airflow/operators/python.py", line 118, in execute_callable
250+
return self.python_callable(*self.op_args, **self.op_kwargs)
251+
Exception: Task failed!
252+
result:
253+
custom:
254+
error:
255+
kind: "Exception"
256+
message: "Task failed!"
257+
stack: |2-
258+
File "/usr/local/lib/python3.11/site-packages/airflow/operators/python.py", line 118, in execute_callable
259+
return self.python_callable(*self.op_args, **self.op_kwargs)
260+
Exception: Task failed!
261+
filename: "taskinstance.py"
262+
level: "ERROR"
263+
lineno: 1058
264+
timestamp: 1744305301220
265+
message: "Task failed!"
266+
status: "error"
267+
tags:
268+
- "source:LOGS_SOURCE"
269+
timestamp: 1744305301220
217270
-
218271
sample: "[2025-02-18 14:46:00,123] {{python_operator.py:105}} INFO - Task started"
219272
tags:

0 commit comments

Comments
 (0)