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]
