ChahatKumar opened a new issue, #74025:
URL: https://github.com/apache/airflow/issues/74025

   ### Under which category would you file this issue?
   
   Airflow Core
   
   ### Apache Airflow version
   
   3.2.2
   
   ### What happened and how to reproduce it?
   
   ### Issue description
   
   Large nested JSON values in `DagRun.conf` are materialized as Python object 
graphs and carried through multiple task-startup boundaries, including database 
deserialization and task workload serialization. This can use several times the 
wire payload size in memory and creates avoidable worker/supervisor memory 
pressure when multiple tasks start concurrently.
   
   An illustrative 12,122,014-byte compact JSON payload containing 11,000 
nested records produced:
   
   | Representation | Size |
   | --- | ---: |
   | Compact JSON | 12,122,014 bytes |
   | Parsed Python object graph | 55,294,806 bytes |
   | Peak allocation while running `json.loads` | 55,296,230 bytes |
   
   These Python figures are a local materialization baseline, not an Airflow 
pod RSS measurement. Airflow may hold additional copies while reading 
`DagRun.conf`, constructing the workload, encoding it, and decoding it in the 
task process.
   
   ### Difference from Airflow 2
   
   Airflow 2's CeleryExecutor sent a small task command and launched the task 
through `airflow tasks run`. It did not send the fully materialized 
task-startup context from the supervisor to a child process through Task SDK 
MessagePack.
   
   Airflow 2 still had to read and decode `DagRun.conf`, so its memory use was 
not zero. Airflow 3 adds another boundary where the context is materialized, 
MessagePack-encoded, transferred, and decoded again before the task body starts.
   
   ### Steps to reproduce
   
   1. Run Airflow 3.2.2 with PostgreSQL and CeleryExecutor.
   2. Create a DAG whose task reads `dag_run.conf`.
   3. Trigger the DAG with a nested `DagRun.conf` payload close to 12 MB.
   4. Start several runs or tasks concurrently.
   5. Measure worker/supervisor RSS while Airflow reads the JSONB value, builds 
the task workload, serializes it, starts the task process, and decodes it.
   6. Compare Airflow 3 with the Airflow 2 Celery task-launch path using the 
same payload and concurrency.
   
   The object-graph amplification can be reproduced independently with:
   
   ```python
   import json
   import tracemalloc
   
   payload = open("large-conf.json", encoding="utf-8").read()
   tracemalloc.start()
   value = json.loads(payload)
   current, peak = tracemalloc.get_traced_memory()
   print(len(payload.encode()), peak)
   ```
   
   
   ### What you think should happen instead?
   
   Airflow should avoid repeatedly materializing and copying large nested 
`DagRun.conf` values during task startup.
   
   Possible approaches include transporting persisted JSON as opaque serialized 
bytes, delaying decoding until the task accesses the value, or otherwise 
avoiding serialization and deserialization of the full Python object graph at 
each boundary.
   
   ### Operating System
   
   _No response_
   
   ### Deployment
   
   Official Apache Airflow Helm Chart
   
   ### Apache Airflow Provider(s)
   
   _No response_
   
   ### Versions of Apache Airflow Providers
   
   apache-airflow-providers-celery
   
   ### Official Helm Chart version
   
   Not Applicable
   
   ### Kubernetes Version
   
   _No response_
   
   ### Helm Chart configuration
   
   Not Applicable
   
   ### Docker Image customizations
   
   Not Applicable
   
   ### Anything else?
   
   Related issues:
   
   - #74023 tracks number representation and Python-type changes when 
PostgreSQL JSONB-backed `DagRun.conf` is materialized. This issue covers memory 
amplification and applies even when all values remain within normal numeric 
ranges.
   - #73712 tracks MessagePack integer overflow during the same task-startup 
transport. This issue is specifically about memory amplification for large 
nested payloads.
   
   ### Are you willing to submit PR?
   
   - [ ] Yes I am willing to submit a PR!
   
   ### Code of Conduct
   
   - [x] I agree to follow this project's [Code of 
Conduct](https://github.com/apache/airflow/blob/main/CODE_OF_CONDUCT.md)
   


-- 
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]

Reply via email to