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

   ## Problem
   
   A Python task gets config-backed defaults such as `[core] 
default_task_retries` and `[operators] default_queue` through the 
`client_defaults` of its serialized payload. A Lang-SDK (TypeScript or Java) 
payload carries no `client_defaults`, so its tasks load with the schema 
defaults instead.
   
   For example, with this config:
   
   ```ini
   [core]
   default_task_retries = 3
   
   [operators]
   default_queue = celery_q
   ```
   
   a task that sets neither `retries` nor `queue` loads as:
   
   - Python task: `retries` 3, `queue` `celery_q`
   - TypeScript or Java task: `retries` 0, `queue` `default`
   
   ## Why it needs runtime changes first
   
   Airflow could fill `client_defaults.tasks` from the config when a payload 
has none, but that is not safe yet. Both runtime serializers drop a task field 
whose value equals its schema default: 
[`applySchemaFields`](https://github.com/apache/airflow/blob/c31888641b1c2c9f3d90b44be1241d7be0b46472/ts-sdk/src/coordinator/serde.ts#L611)
 in the TypeScript SDK on main (#73441), and `serializeTask` in the Java SDK's 
`Serde.kt` ([lines 
139-143](https://github.com/apache/airflow/blob/eb662eac6fcb8c401b4ea45bbda6d0f40a3753bb/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde.kt#L139-L143),
 #71190, not merged yet). So an explicit `retries=0` or `queue="default"` looks 
the same as an unset field, and a server-side fill would override it with the 
config value. A Python task keeps such an explicit value when it comes from the 
operator arguments or `default_args`.
   
   In a prototype of the fill, the stock serializers let explicit default 
values be overridden: a task with `retries=0` and `queue="default"` loaded 3 
and `celery_q`, while the same Python task keeps 0 and `default`. With both 
serializers writing every field the author set, the same fill corrected 164 
values to match Python and broke none (27 test tasks).
   
   Java has a related gap at Dag level. #74041 fills `max_active_tasks`, 
`max_active_runs`, `max_consecutive_failed_dag_runs`, `catchup` and 
`disable_bundle_versioning` from the config, but only when the payload leaves 
them out. The Java serializer always writes them, with the fallbacks 16, 16, 0, 
false and false ([lines 
173-177](https://github.com/apache/airflow/blob/eb662eac6fcb8c401b4ea45bbda6d0f40a3753bb/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde.kt#L173-L177)),
 so that fill never applies to a Java Dag.
   
   ## Plan
   
   In this order:
   
   1. TypeScript serializer: write every task field the author set, even when 
it equals the schema default.
   2. Java serializer: write every config value the author set, and stop 
writing the Dag-level fallbacks (write those five settings only when the Dag 
sets them).
   3. Conformance compare (`scripts/ci/lang_sdk_serialization/compare.py`): 
accept a key the SDK writes when it equals the schema default.
   4. Core: fill `client_defaults.tasks` from the config when the payload has 
none, limited to the config-backed task keys (`owner`, `queue`, `retries`, 
`retry_delay`, `weight_rule`, `execution_timeout`, `email_on_failure`, 
`email_on_retry`).
   5. ADR: update the task encoding section of 
`airflow-core/adr/lang-sdk/0004-dag-parsing.md` to say a runtime writes every 
task field the author set, and Airflow fills the unset config-backed ones from 
its config.
   
   ## Key changes
   
   Trimmed from the prototype diffs.
   
   TypeScript (`ts-sdk/src/coordinator/serde.ts`):
   
   ```diff
    interface FieldRules {
      ...
   +  /** Write a value the author set even when it equals the schema default. 
*/
   +  readonly keepDefaults?: boolean;
    }
   
    const TASK_FIELD_RULES: FieldRules = {
      ...
   +  keepDefaults: true,
    };
   
    function applySchemaFields(data, spec, rules, label) {
      ...
   -    if (value === undefined || value === field.default) continue;
   +    if (value === undefined || (!rules.keepDefaults && value === 
field.default)) continue;
   ```
   
   Kotlin 
(`java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde.kt`):
   
   ```diff
      def.configValues.forEach { (key, value) ->
   -    if (!matchesSchemaDefault(SchemaFields.TASK[key], value)) {
   -      data[key] = unwrapTypeEncoding(serializeValue(value))
   -    }
   +    data[key] = unwrapTypeEncoding(serializeValue(value))
      }
   ```
   
   `matchesSchemaDefault` is then unused and removed. The Dag-level part (write 
the five settings only when the Dag sets them) is not in the prototype.
   
   Python (`DagSerialization.fill_config_defaults` in 
`airflow-core/src/airflow/serialization/serialized_objects.py`, added by 
#74041):
   
   ```diff
                if key not in dag:
                    dag[key] = get(section, option)
   +        if "client_defaults" not in serialized_obj and (
   +            task_defaults := 
OperatorSerialization.generate_client_defaults()
   +        ):
   +            serialized_obj["client_defaults"] = {"tasks": 
dict(task_defaults)}
   ```
   
   The prototype fills everything `generate_client_defaults()` returns, which 
today is only the config-backed keys. The real change should keep only those 
keys.
   
   ## Alternative
   
   The coordinator could pass these options to the runtime through its 
environment, as it already does for `[operators] default_deferrable` 
(`_build_runtime_env()` in 
`task-sdk/src/airflow/sdk/coordinators/_subprocess.py`, #74042), and each 
runtime could fill unset task fields itself. The schema-default skip would then 
stay correct and the server would not need a fill. Not prototyped.
   
   ## Related
   
   - #74041 and its [review 
thread](https://github.com/apache/airflow/pull/74041#discussion_r4159728137)
   - #73441: TypeScript SDK serializer (merged)
   - #71190: Java SDK serializer (draft)
   


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