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]

Reply via email to