uranusjr commented on code in PR #73457:
URL: https://github.com/apache/airflow/pull/73457#discussion_r4154000337


##########
airflow-core/adr/dag-processing/0001-dag-importer-process-model.md:
##########
@@ -0,0 +1,374 @@
+<!--
+ Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements.  See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership.  The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License.  You may obtain a copy of the License at
+
+   http://www.apache.org/licenses/LICENSE-2.0
+
+ Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied.  See the License for the
+ specific language governing permissions and limitations
+ under the License.
+ -->
+
+# ADR-0001: Dag Importer Process Model — Who Owns the Parse Process
+
+## Status
+
+Proposed
+
+## Context
+
+[AIP-85](https://cwiki.apache.org/confluence/x/_Q7OEg) adds 
`AbstractDagImporter`
+(`task-sdk/src/airflow/sdk/importers/`) so the Dag processor can parse sources 
other than Python.
+The interface as it stands says what importing *means* for a format — 
`list_dag_definitions`,
+`import_definition`, `get_source_code` — and says nothing about *where* the 
import runs.
+
+Where it runs is still decided by the Dag processor manager, which forks one
+`DagFileProcessorProcess` per file and runs a Python parse in it. That is the 
wrong shape for both
+directions the AIP opens:
+
+- A format that only has to be **read** — JSON, YAML, a manifest naming Dags — 
pays for a process
+  it does not need, on every parse of every file.
+- A format backed by **another runtime** needs a *different* process — a JVM, 
a Go binary — not a
+  Python fork. Nothing in the current model can express that, which is why Dag 
parsing was cut
+  from AIP-108's scope and [ADR-0004](../lang-sdk/0004-dag-parsing.md) was 
left "retained for when
+  AIP-85 is revisited".
+
+An earlier proposal filled the gap by inserting a layer between the manager 
and the parse process,
+mirroring the executor's worker pool: the manager forks a generic parse 
worker, which then
+dispatches to an importer. This ADR settles the process model without that 
layer, and fixes the
+handful of interface obligations that follow from it.
+
+Terms: a **definition** is the smallest thing an importer can import on its 
own. For Python that is
+a module; for a zip archive it is a member, not the archive. One definition 
may yield several Dags.
+
+## Decision
+
+### 1. The importer owns the process, and nothing sits between it and the 
manager
+
+```
+today
+
+  manager loop
+     │  one queue entry per FILE
+     ▼
+  fork DagFileProcessorProcess              ← always a Python fork, whatever 
the format
+     └── _parse_file → DagBag(file) → every Dag in the file, one sys.modules
+
+decided
+
+  manager loop
+     │  one queue entry per DEFINITION
+     ▼
+  importer.start_import(definition, bundle, context=...) → ImportHandle
+     ├── read it here, already finished        static format — no process at 
all
+     ├── fork_import() → parse child           Python — one child per 
definition
+     └── hand to a runtime it keeps warm       Java / Go — the importer's own 
process
+```
+
+`start_import` is the extension point for *how the importer's language runs*; 
`import_definition`
+remains the single statement of what importing means, and both the inline path 
and the child end up
+back in it, so there is never a second implementation of the import itself.
+
+The default `start_import` imports inline and returns an already-finished 
handle. A static format
+therefore gets no process without asking for one. An importer that executes 
definition-author code
+must override — the manager's own process is not a safe place to run it.
+
+The manager cannot call `import_definition` directly because its loop is 
single-threaded and also
+services sockets, refreshes bundles, harvests results and enforces deadlines; 
`import_definition`
+blocks for as long as user code takes. `ImportHandle` is the shape it 
supervises instead:
+
+```
+ImportHandle                       what the manager does with it
+  is_ready    ──────────────────►  polled every loop pass — must never block
+  result      ──────────────────►  read once ready, then bagged and persisted
+  start_time  ──────────────────►  compared against the parse timeout
+  pid         ──────────────────►  diagnostics and log lines only
+  kill()      ──────────────────►  the deadline passed
+  close()     ──────────────────►  exactly once, after the result is taken
+```
+
+
+The interaction across the manager, importer, and execution boundary follows:
+
+```
+Manager                           Importer                       Worker / 
Subprocess
+   │                                 │                                    │
+   │── start_import(def, context) ──►│                                    │
+   │                                 │── fork_import() ──────────────────►│ 
(spawns child)
+   │◄── ImportHandle ────────────────│                                    │
+   │                                                                      │
+   │── [loop] poll handle.is_ready ──────────────────────────────────────►│ 
(runs import_definition)
+   │                                                                      │
+   │◄── is_ready = True ──────────────────────────────────────────────────│
+   │── read handle.result & close()
+```
+
+The protocol is structural, not a base class: the existing parse subprocess 
already has this shape,
+and an importer wrapping a foreign runtime should not have to inherit from the 
Task SDK to be
+supervised.
+
+#### Alternatives considered
+
+- **A worker-pool layer between manager and importer**, by analogy with the 
executor. Rejected: the
+  analogy does not hold. On the execution side the first fork is a 
*distribution boundary* — the
+  workload may leave the machine. Parsing has no such boundary; the layer 
would only add a process
+  hop and a second place where process policy lives.
+- **A thread pool in the manager.** Rejected: it contradicts fork safety — a 
fork wants no threads
+  running and no locks held — and it would make thread safety a requirement of 
every third-party
+  importer.
+- **Keep the manager forking, and consult the importer only for the work 
inside the child.**
+  Rejected: it cannot express "no process" or "a process of a different kind", 
which is the point.
+- **`async`/`await` on the importer interface.** Rejected: it colours the 
whole interface and drags
+  every importer author into an event loop for what is, for most formats, a 
`read()`.
+
+### 2. The unit of parse work is a Dag definition, not a file
+
+```
+bundle/
+  daily.py               → 1 definition    daily.py              (1 Dag)
+  etl.py                 → 1 definition    etl.py                (3 Dags)
+  archive.zip
+    ├── a.py             → 1 definition    archive.zip/a.py
+    └── b.py             → 1 definition    archive.zip/b.py
+                           ^^ previously ONE entry for the whole archive, with 
both
+                              members imported into a single shared sys.modules
+  catalog.yaml           → N definitions   one per Dag the manifest names

Review Comment:
   Yeah it should be 1 def, similar to one Python file.



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