Repository navigation
Resolve OOM when reading large logs in webserver #45079
Description
Activity
- addedkind:featureFeature RequestsFeature Requestsneeds-triagelabel for new issues that we didn't triage yetlabel for new issues that we didn't triage yet
on Dec 19, 2024 Thanks for opening your first issue here! Be sure to follow the issue template! If you are willing to raise PR to address this issue please do so, no need to wait for approval.
Yes. That's exactly how I envisioned solving this problem. @dstandish ?
Reacted by Jason(Zhe-You) LiuFYI. Breaking changes to FileTaskHandler is not a problem - we can work out back-compatibility or simply break it for Airflow 3 - this is not a big deal, since this is only a deployment configuration and does not require DAG adaptations.
Reacted by Jason(Zhe-You) Liu- removedneeds-triagelabel for new issues that we didn't triage yetlabel for new issues that we didn't triage yet
on Dec 19, 2024 Hi @potiuk,
Would it be okay if I treat this issue as an umbrella issue to track other TODO tasks while refactoring each provider? Or would it be more preferable to refactor
FileTaskHandlerand all providers in a single PR in this case? Thanks !Sure. It can be separate set of PRs and that issue can remain "umbrella" - you do not need to have more issues. PRs are enough
Reacted by Jason(Zhe-You) LiuYes. That's exactly how I envisioned solving this problem. @dstandish ?
IIRC this should be fine when task done but may present challenges when task is in flight because at any moment the location of the logs may shift eg from worker to remote storage etc
IIRC this should be fine when task done but may present challenges when task is in flight because at any moment the location of the logs may shift eg from worker to remote storage etc
Is it not the same case now?
Related issue #31105
Reacted by Jason(Zhe-You) LiuYes. That's exactly how I envisioned solving this problem. @dstandish ?
IIRC this should be fine when task done but may present challenges when task is in flight because at any moment the location of the logs may shift eg from worker to remote storage etc
Taking
S3TaskHandleras an example, it requires additional refactoring and might need aread_streammethod added toS3Hookthat returns a generator-based result:
https://github.com/apache/airflow/blob/main/providers/src/airflow/providers/amazon/aws/log/s3_task_handler.py#L136-L192From my perspective, for the
s3_writecase, I would download the old log as temporary file and append the new log stream into a temporary file, and use theupload_filemethod to upload the file to prevent memory starvation and remain the same result.Taking S3TaskHandler as an example, it requires additional refactoring and might need a read_stream method added to S3Hook that returns a generator-based result:
https://github.com/apache/airflow/blob/main/providers/src/airflow/providers/amazon/aws/log/s3_task_handler.py#L136-L192From my perspective, for the s3_write case, I would download the old log as temporary file and append the new log stream into a temporary file, and use the upload_file method to upload the file to prevent memory starvation and remain the same result.
Yep. There will be dga cases like that. And yes the proposed method is good.
Reacted by Jason(Zhe-You) LiuResolve OOM When Reading Large Logs in Webserver #49470 and its backport [v3-0-test] Resolve OOM When Reading Large Logs in Webserver #53167 have been merged, so this issue can now be closed as completed.
The patch will be included in the next Airflow release,
3.0.4.
Description
Related context: #44753 (comment)
TL;DR
After conducting some research and implementing a POC, I would like to propose a potential solution. However, this solution requires changes to the
airflow.utils.log.file_task_handler.FileTaskHandler. If the solution is accepted, it will necessitate modifications to 10 providers that extend theFileTaskHandlerclass.Main Concept for Refactoring
The proposed solution focuses on:
The POC for this refactoring shows a 90% reduction in memory usage with similar processing times!
Experiment Details
Main Root Causes of OOM
_interleave_logsFunction inairflow.utils.log.file_task_handlerrecordslist.recordslist._readMethod inairflow.utils.log.file_task_handler.FileTaskHandler_read:These methods read the entire log content and return it as a string instead of a generator:
_read_from_local_read_from_logs_server_read_remote_logs(Implemented by providers)Proposed Refactoring Solution
The main concept includes:
heapqwith streams of logs.Breaking Changes in This Solution
Interface of the
readMethod inFileTaskHandler:Interfaces of
read_log_chunksandread_log_streaminTaskLogReader:Methods That Use
_read_read_from_local_read_from_logs_server_read_remote_logs( there are 10 providers implement this method )Experimental Environment:
830 MB, about8670000linesBenchmark Metrics
Original Implementation:
POC (Refactored Implementation):
Summary
Feel free to share any feedback! I believe we should have more discussions before adopting this solution, as it involves breaking changes to the
FileTaskHandlerinterface and requires refactoring in 10 providers as well.Related issues
#44753
TODO Tasks
FileTaskHandler_read_remote_logsmethod_readwith_read_remote_logsin Rework remote task log handling for the structlog era. #48491 )_readmethodAmazon AWS Cloud Watchget_task_logCeleryKubernetesExecutorKubernetesExecutorLocalKubernetesExecutorAre you willing to submit a PR?
Code of Conduct