-
Notifications
You must be signed in to change notification settings - Fork 4
/
loki_task_handler.py
202 lines (156 loc) · 6.07 KB
/
loki_task_handler.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
"""Loki logging handler for tasks"""
import time
from typing import Any, Dict, Optional, Tuple, List
import typing
import gzip
if typing.TYPE_CHECKING:
from airflow.models import TaskInstance
from typing import Optional, Tuple
import logging
logging.raiseExceptions = True
import time
from typing import Optional
import json
from datetime import timedelta
from airflow.utils.log.file_task_handler import FileTaskHandler
from airflow.utils.log.logging_mixin import LoggingMixin
from grafana_loki_provider.hooks.loki import LokiHook
from airflow.compat.functools import cached_property
from airflow.configuration import conf
import os
BasicAuth = Optional[Tuple[str, str]]
DEFAULT_LOGGER_NAME = "airflow"
import json
import logging
import typing
class LokiTaskHandler(FileTaskHandler, LoggingMixin):
def __init__(
self,
base_log_folder,
name,
filename_template: Optional[str] = None,
enable_gzip=True,
):
super().__init__(base_log_folder, filename_template)
self.name: str = name
self.handler: Optional[logging.FileHandler] = None
self.log_relative_path = ""
self.closed = False
self.upload_on_close = True
self.enable_gzip = enable_gzip
self.labels: Dict[str, str] = {}
self.extras: Dict[str, Any] = {}
@cached_property
def hook(self) -> LokiHook:
"""Returns LokiHook"""
remote_conn_id = str(conf.get("logging", "REMOTE_LOG_CONN_ID"))
from grafana_loki_provider.hooks.loki import LokiHook
return LokiHook(loki_conn_id=remote_conn_id)
def get_extras(self, ti, try_number=None) -> Dict[str, Any]:
return dict(
run_id=getattr(ti, "run_id", ""),
try_number=try_number if try_number != None else ti.try_number,
map_index=getattr(ti, "map_index", ""),
)
def get_labels(self, ti) -> Dict[str, str]:
return {"dag_id": ti.dag_id, "task_id": ti.task_id}
def set_context(self, task_instance: "TaskInstance") -> None:
super().set_context(task_instance)
ti = task_instance
self.log_relative_path = self._render_filename(ti, ti.try_number)
self.upload_on_close = not ti.raw
# Clear the file first so that duplicate data is not uploaded
# when re-using the same path (e.g. with rescheduled sensors)
if self.upload_on_close:
if self.handler:
with open(self.handler.baseFilename, "w"):
pass
self.labels = self.get_labels(ti)
self.extras = self.get_extras(ti)
def _get_task_query(self, ti, try_number, metadata) -> str:
run_id = getattr(ti, "run_id", "")
map_index = getattr(ti, "map_index", "")
query = """ {{dag_id="{dag_id}",task_id="{task_id}"}}
| json try_number="try_number",map_index="map_index",run_id="run_id"
| try_number="{try_number}" and
map_index="{map_index}" and
run_id="{run_id}"
| __error__!="JSONParserErr"
""".format(
try_number=try_number,
map_index=map_index,
run_id=run_id,
dag_id=ti.dag_id,
task_id=ti.task_id,
)
return query
def _read(
self, ti, try_number: int, metadata: Optional[str] = None
) -> Tuple[str, Dict[str, bool]]:
query = self._get_task_query(ti, try_number, metadata)
start = ti.start_date - timedelta(days=15)
#if the task is running or queued, the task will not have end_date, in that
# case, we will use a resonable internal of 5 days
end_date = ti.end_date or ti.start_date + timedelta(days=5)
end = end_date + timedelta(hours=1)
params = {
"query": query,
"start": start.isoformat(),
"end": end.isoformat(),
"limit":5000,
"direction": "forward",
}
self.log.info(f"loki log query params {params}")
data = self.hook.query_range(params)
lines = []
if "data" in data and "result" in data["data"]:
for i in data["data"]["result"]:
for v in i["values"]:
try:
msg = v[1]
line = json.loads(msg)["line"]
lines.append(line)
except Exception as e:
self.log.exception(e)
pass
if lines:
log_lines = "".join(lines)
return log_lines, {"end_of_log": True}
else:
return super()._read(ti, try_number, metadata)
def close(self):
"""Close and upload local log file to remote storage Loki."""
if self.closed:
return
super().close()
if not self.upload_on_close:
return
local_loc = os.path.join(self.local_base, self.log_relative_path)
if os.path.exists(local_loc):
# read log and remove old logs to get just the latest additions
with open(local_loc) as logfile:
log = logfile.readlines()
self.loki_write(log)
# Mark closed so we don't double write if close is called twice
self.closed = True
def build_payload(self, log: List[str], labels, extras) -> dict:
"""Build JSON payload with a log entry."""
ns = 1e9
lines = []
for line in log:
ts = str(int(time.time() * ns))
line = {**{"line": line}, ** extras }
line = json.dumps(line)
lines.append([ts, line])
stream = {
"stream": labels,
"values": lines,
}
return {"streams": [stream]}
def loki_write(self, log):
payload = self.build_payload(log, self.labels, self.extras)
headers = {"Content-Type": "application/json"}
if self.enable_gzip:
payload = gzip.compress(json.dumps(payload).encode("utf-8"))
headers["Content-Encoding"] = "gzip"
self.hook.push_log(payload=payload, headers=headers)