-
Notifications
You must be signed in to change notification settings - Fork 1.3k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Fix celery docker jobs that involve writing directly to command-line …
…output (#7665) Summary: This brings in the celery docker executor to match the logic in the celery k8s executor for turning process logs into structured events. Before, logging in an op or resources would lead to a cryptic JSON decode error. Test Plan: BK
- Loading branch information
Showing
7 changed files
with
78 additions
and
57 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,47 @@ | ||
from json import JSONDecodeError | ||
|
||
from dagster import check | ||
from dagster.serdes import deserialize_json_to_dagster_namedtuple | ||
|
||
|
||
def filter_dagster_events_from_cli_logs(log_lines): | ||
""" | ||
Filters the raw log lines from a dagster-cli invocation to return only the lines containing json. | ||
- Log lines don't necessarily come back in order | ||
- Something else might log JSON | ||
- Docker appears to silently split very long log lines -- this is undocumented behavior | ||
TODO: replace with reading event logs from the DB | ||
""" | ||
check.list_param(log_lines, "log_lines", str) | ||
|
||
coalesced_lines = [] | ||
buffer = [] | ||
in_split_line = False | ||
for line in log_lines: | ||
line = line.strip() | ||
if not in_split_line and line.startswith("{"): | ||
if line.endswith("}"): | ||
coalesced_lines.append(line) | ||
else: | ||
buffer.append(line) | ||
in_split_line = True | ||
elif in_split_line: | ||
buffer.append(line) | ||
if line.endswith("}"): # Note: hack, this may not have been the end of the full object | ||
coalesced_lines.append("".join(buffer)) | ||
buffer = [] | ||
in_split_line = False | ||
|
||
events = [] | ||
for line in coalesced_lines: | ||
try: | ||
events.append(deserialize_json_to_dagster_namedtuple(line)) | ||
except JSONDecodeError: | ||
pass | ||
except check.CheckError: | ||
pass | ||
|
||
return events |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters