li3zhi4 opened a new pull request, #11633: URL: https://github.com/apache/seatunnel/pull/11633
### Purpose of this pull request `schema.columns[].defaultValue` is parsed and stored on `Column` (`ReadonlyConfigParser` → `PhysicalColumn.of(...)`), and `Column.getDefaultValue()` is part of the public API — but the **JSON format deserializer never applies it**: `JsonToRowConverters.convertField()` returns `null` for missing/`null` fields without consulting the column metadata, because `JsonDeserializationSchema` only passes `SeaTunnelRowType` (field names + types) and never hands over `Column[]`. Result: every connector using the JSON format (Kafka, Pulsar, HTTP, MongoDB, Elasticsearch, ...) silently produces `null` where the user configured a default value — a data-integrity bug. The related Debezium-path fix ([PR #7950](https://github.com/apache/seatunnel/commit/3b432125a)) did not touch this code path. ### What changed #### 1. `JsonToRowConverters.java` - Add a `Column[] columns` field and a new constructor `JsonToRowConverters(boolean failOnMissingField, boolean ignoreParseErrors, Column[] columns)`; the existing 2-arg constructor is kept and delegates with `columns = null` (fully backward compatible). - Add a `convertField(JsonToObjectConverter, String, JsonNode, int fieldIndex)` overload: when the field is **missing or `null`**, it consults `columns[fieldIndex].getDefaultValue()` and converts it to the field type via the field converter, falling back to the original `failOnMissingField`/`null` behavior when no default value is configured. Normalizing through the field converter also handles defaults whose raw type differs from the column type (e.g. HOCON config parses `0.0` as `Integer 0`, which must become a `Double` for a double column). ```java private Object convertField( JsonToObjectConverter fieldConverter, String fieldName, JsonNode field, int fieldIndex) { if (field == null || field.isNull()) { if (columns != null && fieldIndex < columns.length) { Object defaultValue = columns[fieldIndex].getDefaultValue(); if (defaultValue != null) { return fieldConverter.convert(JsonUtils.toJsonNode(defaultValue), fieldName); } } if (failOnMissingField) { throw new IllegalArgumentException( String.format("Could not find field with name %s .", fieldName)); } else { return null; } } else { return fieldConverter.convert(field, fieldName); } } ``` - `createRowConverter` calls the new overload with the field index. #### 2. `JsonDeserializationSchema.java` - Extract `Column[]` from `catalogTable.getTableSchema().getColumns()` (with null-safety on both the table schema and the column list) and pass it to the new `JsonToRowConverters` constructor. The `SeaTunnelRowType`-only constructor keeps passing `null` (no metadata available → old behavior). #### 3. New unit test `JsonDefaultValueTest.java` Covers: field missing → defaultValue applied; explicit `null` → defaultValue applied; no defaultValue configured → `null` (backward compatibility); multiple data types (string / int / double / boolean / long / float / decimal / date / timestamp); multiple double number formats (scientific notation, negative values, non-zero decimals, integer-valued doubles, string-typed numeric defaults); and Integer defaultValue normalized to a double field type. #### 4. New e2e test `KafkaJsonDefaultValueIT.java` (connector-kafka-e2e) Kafka source (`format = json` + `schema.columns[].defaultValue`) → Kafka sink (json). Produces 3 messages (field missing / explicit null / real values incl. scientific notation `2e3`) against a 10-column schema (string / int / double / boolean / bigint / float / decimal / date / timestamp), then consumes the sink topic and asserts the default values are applied with the correct types. Restricted to the Zeta engine (`@DisabledOnContainer` for FLINK/SPARK) since the fix lives in the engine-agnostic `seatunnel-format-json` module. ### How was this patch tested? - Unit tests: `JsonDefaultValueTest` 8/8 pass; full `seatunnel-format-json` module 49/49 pass (`BUILD SUCCESS`). - e2e: `KafkaJsonDefaultValueIT` 1/1 pass on the Zeta engine (`Tests run: 1, Failures: 0, Errors: 0` + `BUILD SUCCESS`). - `./mvnw spotless:apply` clean. ### Check list - [x] Code changes are complete and verified locally - [x] `mvn spotless:apply` passes - [x] Unit tests pass (49/49) - [x] e2e test added and passing (1/1, Zeta engine) - [x] Issue(s) referenced: closes #11632 - [ ] Docs updated if needed (schema defaultValue doc already exists upstream, no doc change needed) -- 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]
