This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12057-dc7c3d0ab5e4eef98987e1f5840f505ab261448c in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit b0dd4945b1da8423a481b053082bd5fd11b9533f Author: 雷炯 <[email protected]> AuthorDate: Sun Oct 4 23:00:07 2026 +0000 [Feature][Transform-V2] Support connector-declared metadata fields in… (#12057) Co-authored-by: Cursor <[email protected]> --- docs/en/transforms/metadata.md | 29 ++- docs/zh/transforms/metadata.md | 29 ++- .../transform/metadata/MetadataTransform.java | 38 ++- .../MetadataMultiCatalogSchemaChangeTest.java | 107 ++++++-- .../transform/metadata/MetadataTransformTest.java | 274 +++++++++++++++++++++ 5 files changed, 456 insertions(+), 21 deletions(-) diff --git a/docs/en/transforms/metadata.md b/docs/en/transforms/metadata.md index f1ca25eead..d7df6c91f3 100644 --- a/docs/en/transforms/metadata.md +++ b/docs/en/transforms/metadata.md @@ -33,6 +33,31 @@ The Metadata transform plugin is used to extract metadata information from data | Gtid | string | Global Transaction ID (`server_uuid:transaction_id`). `null` when GTID is disabled or for snapshot rows. | MySQL-CDC only | | Partition | string | Partition information of the data, multiple partition fields separated by commas | Connectors supporting partitions | +## Connector-Declared Metadata Fields + +Connectors can expose source-specific metadata through `CatalogTable.MetadataSchema` without registering every key in the global metadata registry. The `Metadata` transform accepts a logical key when **either**: + +1. It is a globally supported metadata key (see the tables in this document), or +2. It is explicitly declared by the input table's `MetadataSchema` + +For connector-declared keys, the output physical column keeps the declared data type, nullability, length, default value, and comment. Runtime values are read from `SeaTunnelRow.options`. + +This does **not** allow arbitrary row-option keys. A key present only in `SeaTunnelRow.options`, but absent from both the global registry and `MetadataSchema`, is rejected. Matching is case-sensitive. The transform does not read physical columns that happen to use the same name. + +```hocon +transform { + Metadata { + plugin_input = "source_rows" + plugin_output = "rows_with_source_metadata" + metadata_fields { + KafkaOffset = kafka_offset + } + } +} +``` + +The upstream source or transform must declare `KafkaOffset` in `CatalogTable.metadataSchema` and write the corresponding value into `SeaTunnelRow.options`. The `Metadata` transform does not invent connector-specific metadata by itself. + ## Knowledge Sync Metadata Fields Knowledge Sync pipelines can use the following logical metadata keys to carry document and chunk identity. These keys become physical fields only after they are explicitly projected by the `Metadata` transform. @@ -54,7 +79,7 @@ The `Metadata` transform does not generate Knowledge Sync metadata by itself. Th ### Important Notes -1. **Metadata field names are case-sensitive**: Configuration must strictly follow the Key names in the table above (e.g., `Database`, `Table`, `RowKind`, etc.) +1. **Metadata field names are case-sensitive**: Configuration must strictly follow the Key names in the tables above (e.g., `Database`, `Table`, `RowKind`) or the exact names declared in the input `MetadataSchema`. 2. **Time fields**: `Delay` and `SourceTimestamp` are only available for CDC connectors. `EventTime` is also provided by the Kafka source via `ConsumerRecord.timestamp` when available. 3. **Kafka event time**: The Kafka source writes `ConsumerRecord.timestamp` (milliseconds) into `EventTime` when it is non-negative, so you can surface it with the `Metadata` transform. 4. **Binlog/GTID fields**: `BinlogFile`, `BinlogPos`, `BinlogRow`, and `Gtid` are MySQL-CDC specific. For `startup.mode = initial`, snapshot rows return `null` for all four fields. @@ -137,7 +162,7 @@ metadata_fields { ``` **Notes:** -- The left side must be a supported metadata Key (see table above), and is strictly case-sensitive +- The left side must be a globally supported metadata Key (see the tables above) or a key declared in the input table `MetadataSchema`, and is strictly case-sensitive - The right side is a custom output field name, which cannot duplicate existing field names - You can select only the metadata fields you need, not all of them must be configured diff --git a/docs/zh/transforms/metadata.md b/docs/zh/transforms/metadata.md index ee3520998e..c4eab39527 100644 --- a/docs/zh/transforms/metadata.md +++ b/docs/zh/transforms/metadata.md @@ -33,6 +33,31 @@ Metadata 转换插件用于将数据行中的元数据信息提取为普通字 | Gtid | string | 全局事务 ID(格式:`server_uuid:transaction_id`)。GTID 未启用或快照行时返回 `null`。 | 仅 MySQL-CDC | | Partition | string | 数据所属的分区信息,多个分区字段使用逗号分隔 | 支持分区的连接器 | +## 连接器声明的元数据字段 + +连接器可以通过 `CatalogTable.MetadataSchema` 暴露各自的元数据,而不必把每个 Key 都注册到全局元数据表中。`Metadata` 转换会接受满足以下任一条件的逻辑 Key: + +1. 本文档表格中已列出的全局元数据 Key +2. 输入表 `MetadataSchema` 中显式声明的 Key + +对于连接器声明的 Key,输出物理列会保留声明的数据类型、可空性、长度、默认值和注释。运行时的值从 `SeaTunnelRow.options` 读取。 + +这**不会**放开任意 row option。一个 Key 如果只出现在 `SeaTunnelRow.options` 中,但既不在全局注册表、也不在 `MetadataSchema` 中,仍然会被拒绝。匹配区分大小写。该转换不会读取碰巧同名的物理列。 + +```hocon +transform { + Metadata { + plugin_input = "source_rows" + plugin_output = "rows_with_source_metadata" + metadata_fields { + KafkaOffset = kafka_offset + } + } +} +``` + +上游 Source 或 Transform 必须先在 `CatalogTable.metadataSchema` 中声明 `KafkaOffset`,并把对应值写入 `SeaTunnelRow.options`。`Metadata` 转换本身不会生成连接器特有的元数据。 + ## Knowledge Sync 元数据字段 Knowledge Sync 流程可以使用下面的逻辑元数据 Key 携带文档和 chunk 身份信息。这些 Key 只有通过 `Metadata` 转换显式投影后,才会成为真实的物理列。 @@ -54,7 +79,7 @@ Knowledge Sync 流程可以使用下面的逻辑元数据 Key 携带文档和 ch ### 重要说明 -1. **元数据字段区分大小写**:配置时必须严格按照上表中的 Key 名称(如 `Database`、`Table`、`RowKind` 等)。 +1. **元数据字段区分大小写**:配置时必须严格按照上表中的 Key 名称(如 `Database`、`Table`、`RowKind`),或输入表 `MetadataSchema` 中声明的精确名称。 2. **时间相关字段**:`Delay` 和 `SourceTimestamp` 仅在 CDC 连接器有效。`EventTime` 也会在 Kafka 源中使用 `ConsumerRecord.timestamp`(毫秒,非负时)写入。 3. **Kafka 事件时间**:Kafka 源会在 `ConsumerRecord.timestamp` 非负时写入 `EventTime`,可通过 Metadata 转换将其暴露为普通字段。 4. **Binlog/GTID 字段**:`BinlogFile`、`BinlogPos`、`BinlogRow`、`Gtid` 仅适用于 MySQL-CDC。使用 `startup.mode = initial` 时,快照行的这四个字段均为 `null`。 @@ -137,7 +162,7 @@ metadata_fields { ``` **注意事项:** -- 左侧必须是支持的元数据 Key(见上表),且严格区分大小写 +- 左侧必须是全局支持的元数据 Key(见上表),或输入表 `MetadataSchema` 中声明的 Key,且严格区分大小写 - 右侧是自定义的输出字段名,不能与原有字段重名 - 可以只选择需要的元数据字段,不必全部配置 diff --git a/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/metadata/MetadataTransform.java b/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/metadata/MetadataTransform.java index bb7b8cb411..2765481f66 100644 --- a/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/metadata/MetadataTransform.java +++ b/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/metadata/MetadataTransform.java @@ -41,6 +41,14 @@ import java.util.Map; import static org.apache.seatunnel.api.table.type.MetadataUtil.isMetadataField; +/** + * Projects logical row metadata into physical columns. + * + * <p>A {@code metadata_fields} key is accepted when it is globally registered in {@link + * MetadataUtil} or explicitly declared by the input table's {@link MetadataSchema}. Keys that exist + * only in {@code SeaTunnelRow.options} are rejected so that {@link MetadataSchema} remains the type + * and trust boundary for connector-specific fields. + */ public class MetadataTransform extends MultipleFieldOutputTransform { private List<String> fieldNames; @@ -52,13 +60,23 @@ public class MetadataTransform extends MultipleFieldOutputTransform { initOutputFields(inputCatalogTable, config.get(MetadataTransformConfig.METADATA_FIELDS)); } + /** + * Validates {@code metadata_fields} mappings and records the projection order. + * + * <p>Connector-specific keys must be declared in the input {@link MetadataSchema}. Matching is + * case-sensitive. Duplicate physical output names are still rejected. + * + * @param inputCatalogTable upstream catalog table that supplies the metadata schema + * @param fields mapping from logical metadata key to physical output name + */ private void initOutputFields(CatalogTable inputCatalogTable, Map<String, String> fields) { List<String> sourceTableFiledNames = Arrays.asList(inputCatalogTable.getTableSchema().getFieldNames()); + this.metadataSchema = inputCatalogTable.getMetadataSchema(); List<String> fieldNames = new ArrayList<>(); for (Map.Entry<String, String> field : fields.entrySet()) { String srcField = field.getKey(); - if (!isMetadataField(srcField)) { + if (!isProjectableMetadataField(srcField)) { throw TransformCommonError.cannotFindMetadataFieldError(getPluginName(), srcField); } String targetField = field.getValue(); @@ -68,10 +86,26 @@ public class MetadataTransform extends MultipleFieldOutputTransform { fieldNames.add(field.getKey()); } this.fieldNames = fieldNames; - this.metadataSchema = inputCatalogTable.getMetadataSchema(); this.metadataFieldMapping = fields; } + /** + * Returns whether {@code fieldName} can be projected by this transform. + * + * <p>Globally registered keys keep their existing behavior. Connector-specific keys are allowed + * only when the input {@link MetadataSchema} declares them. Physical columns that happen to use + * the same name are not treated as metadata. + * + * @param fieldName logical metadata key from {@code metadata_fields} + * @return {@code true} if the key is globally registered or schema-declared + */ + private boolean isProjectableMetadataField(String fieldName) { + if (isMetadataField(fieldName)) { + return true; + } + return metadataSchema != null && metadataSchema.contains(fieldName); + } + @Override public String getPluginName() { return MetadataTransformConfig.PLUGIN_NAME; diff --git a/seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/metadata/MetadataMultiCatalogSchemaChangeTest.java b/seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/metadata/MetadataMultiCatalogSchemaChangeTest.java index c4833a603e..4efd98364b 100644 --- a/seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/metadata/MetadataMultiCatalogSchemaChangeTest.java +++ b/seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/metadata/MetadataMultiCatalogSchemaChangeTest.java @@ -44,17 +44,16 @@ import java.util.List; import java.util.Map; /** - * Reproduces upstream Bug 1: schema-change events are silently dropped by {@code - * AbstractMultiCatalogTransform} subclasses, leaving the inner per-table transforms with stale - * {@code inputCatalogTable} after ALTER. Result: post-ALTER rows lose new column values. + * Regression coverage for {@link MetadataMultiCatalogTransform} under live schema-change events. * - * <p>This test calls the OUTER wrapper ({@code MetadataMultiCatalogTransform}) — the same instance - * the SeaTunnel engine constructs via the factory and feeds via {@code TransformFlowLifeCycle}. - * Earlier {@code TransformChainLiveAlterTest} called the inner transform directly, missing the bug. + * <p>Exercises the outer multi-catalog wrapper (the instance the engine builds via the factory and + * feeds through {@code TransformFlowLifeCycle}), verifying that {@code mapSchemaChangeEvent} + * dispatches ALTER events to inner per-table transforms so post-ALTER rows keep the expected + * arity/shape. Also covers connector-declared metadata fields surviving ALTER through the wrapper. * - * <p>Without the fix in {@code AbstractMultiCatalogTransform.mapSchemaChangeEvent}, the wrapper's - * default no-op returns the event without dispatching to inner transforms, so this test FAILS at - * the post-ALTER arity assertion. + * <p>Note: dispatch of schema-change events through {@code AbstractMultiCatalogTransform} is + * pre-existing behavior (not introduced by the connector-declared metadata change); the first case + * below is a general regression guard, while the second case exercises the new Metadata feature. */ public class MetadataMultiCatalogSchemaChangeTest { @@ -150,9 +149,8 @@ public class MetadataMultiCatalogSchemaChangeTest { null, null))); - // This is the exact call TransformFlowLifeCycle.received makes on the outer wrapper. - // Without the fix: wrapper's default no-op returns event unchanged; inner transformMap - // entries never see ALTER; their inputCatalogTable stays at 2 cols. + // Same call path as TransformFlowLifeCycle.received on the outer wrapper: the event must + // be dispatched to the inner MetadataTransform so its catalog/schema stay in sync. wrapper.mapSchemaChangeEvent(alter); // Post-ALTER row: 4 base cols (id, name, discount_pct, is_featured) @@ -167,9 +165,8 @@ public class MetadataMultiCatalogSchemaChangeTest { Assertions.assertEquals( 6, postOut.getArity(), - "post-ALTER MUST be arity 6 (4 base + 2 metadata). If 4, the wrapper " - + "swallowed the schema change without notifying inner MetadataTransform — " - + "exactly Bug 1."); + "post-ALTER MUST be arity 6 (4 base + 2 metadata). If 4, the wrapper did not" + + " propagate the schema change to the inner MetadataTransform."); Assertions.assertEquals(2L, postOut.getField(0)); Assertions.assertEquals("Premium A", postOut.getField(1)); Assertions.assertEquals( @@ -181,4 +178,84 @@ public class MetadataMultiCatalogSchemaChangeTest { postOut.getField(3), "is_featured must survive the wrapper after live ALTER"); } + + @Test + void multiCatalogWrapperProjectsConnectorDeclaredMetadataAfterSchemaChange() { + List<Column> metadata = new ArrayList<>(); + metadata.add( + MetadataColumn.of( + "KafkaOffset", + BasicType.LONG_TYPE, + (Long) null, + true, + null, + "Kafka record offset")); + CatalogTable baseTable = + CatalogTable.of( + TableIdentifier.of("catalog", TBL), + TableSchema.builder() + .column( + PhysicalColumn.of( + "id", + BasicType.LONG_TYPE, + (Long) null, + false, + null, + null)) + .column( + PhysicalColumn.of( + "name", + BasicType.STRING_TYPE, + (Long) null, + true, + null, + null)) + .build(), + new HashMap<>(), + new ArrayList<>(), + "comment", + "test", + MetadataSchema.builder().columns(metadata).build()); + + Map<String, String> metaMapping = new LinkedHashMap<>(); + metaMapping.put("KafkaOffset", "kafka_offset"); + Map<String, Object> cfg = new HashMap<>(); + cfg.put("metadata_fields", metaMapping); + MetadataMultiCatalogTransform wrapper = + new MetadataMultiCatalogTransform( + Collections.singletonList(baseTable), ReadonlyConfig.fromMap(cfg)); + + SeaTunnelRow preRow = new SeaTunnelRow(new Object[] {1L, "Widget A"}); + preRow.setTableId(TBL.getFullName()); + preRow.getOptions().put("KafkaOffset", 42L); + SeaTunnelRow preOut = wrapper.map(preRow); + Assertions.assertEquals(3, preOut.getArity(), "pre-ALTER: 2 base + 1 custom metadata"); + Assertions.assertEquals(42L, preOut.getField(2)); + + TableIdentifier tid = baseTable.getTableId(); + AlterTableColumnsEvent alter = + new AlterTableColumnsEvent(tid) + .addEvent( + AlterTableAddColumnEvent.add( + tid, + PhysicalColumn.of( + "discount_pct", + BasicType.DOUBLE_TYPE, + (Long) null, + true, + null, + null))); + wrapper.mapSchemaChangeEvent(alter); + + SeaTunnelRow postRow = new SeaTunnelRow(new Object[] {2L, "Premium A", 10.00d}); + postRow.setTableId(TBL.getFullName()); + postRow.getOptions().put("KafkaOffset", 99L); + SeaTunnelRow postOut = wrapper.map(postRow); + + Assertions.assertEquals(4, postOut.getArity(), "post-ALTER: 3 base + 1 custom metadata"); + Assertions.assertEquals(2L, postOut.getField(0)); + Assertions.assertEquals("Premium A", postOut.getField(1)); + Assertions.assertEquals(10.00d, postOut.getField(2)); + Assertions.assertEquals(99L, postOut.getField(3)); + } } diff --git a/seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/metadata/MetadataTransformTest.java b/seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/metadata/MetadataTransformTest.java index 2e91b3f4c8..0bf461d892 100644 --- a/seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/metadata/MetadataTransformTest.java +++ b/seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/metadata/MetadataTransformTest.java @@ -43,6 +43,7 @@ import java.time.LocalDateTime; import java.time.ZoneOffset; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collections; import java.util.HashMap; import java.util.LinkedHashMap; import java.util.List; @@ -280,6 +281,251 @@ public class MetadataTransformTest { Assertions.assertEquals("chunk_hash_value", output.getField(16)); } + @Test + void shouldProjectConnectorDeclaredTypedMetadataField() { + CatalogTable table = + catalogTableWithCustomMetadata( + Collections.singletonList( + MetadataColumn.of( + "KafkaOffset", + BasicType.LONG_TYPE, + 19L, + true, + null, + "Kafka record offset"))); + MetadataTransform transform = + new MetadataTransform( + ReadonlyConfig.fromMap(metadataFieldsConfig("KafkaOffset", "kafka_offset")), + table); + transform.initRowContainerGenerator(); + + Column[] columns = transform.getOutputColumns(); + Assertions.assertEquals(1, columns.length); + Assertions.assertEquals("kafka_offset", columns[0].getName()); + Assertions.assertEquals(BasicType.LONG_TYPE, columns[0].getDataType()); + Assertions.assertTrue(columns[0].isNullable()); + Assertions.assertEquals(19L, columns[0].getColumnLength()); + Assertions.assertNull(columns[0].getDefaultValue()); + Assertions.assertEquals("Kafka record offset", columns[0].getComment()); + Assertions.assertInstanceOf(PhysicalColumn.class, columns[0]); + + SeaTunnelRow input = new SeaTunnelRow(new Object[] {"payload"}); + input.getOptions().put("KafkaOffset", 42L); + + SeaTunnelRow output = transform.map(input); + Assertions.assertEquals(2, output.getArity()); + Assertions.assertEquals("payload", output.getField(0)); + Assertions.assertEquals(42L, output.getField(1)); + } + + @Test + void shouldRejectUndeclaredConnectorMetadataKey() { + CatalogTable table = catalogTableWithCustomMetadata(Collections.emptyList()); + Map<String, Object> config = metadataFieldsConfig("KafkaOffset", "kafka_offset"); + + TransformException exception = + Assertions.assertThrows( + TransformException.class, + () -> new MetadataTransform(ReadonlyConfig.fromMap(config), table)); + Assertions.assertTrue(exception.getMessage().contains("KafkaOffset")); + } + + @Test + void shouldRejectConnectorMetadataOutputNameCollision() { + CatalogTable table = + catalogTableWithCustomMetadata( + Collections.singletonList( + MetadataColumn.of( + "KafkaOffset", + BasicType.LONG_TYPE, + (Long) null, + true, + null, + "Kafka record offset"))); + + TransformException exception = + Assertions.assertThrows( + TransformException.class, + () -> + new MetadataTransform( + ReadonlyConfig.fromMap( + metadataFieldsConfig("KafkaOffset", "payload")), + table)); + Assertions.assertTrue(exception.getMessage().contains("KafkaOffset")); + } + + @Test + void shouldRejectPhysicalColumnLookalikeThatIsNotMetadata() { + CatalogTable table = + CatalogTable.of( + TableIdentifier.of("catalog", TablePath.DEFAULT), + TableSchema.builder() + .column( + PhysicalColumn.of( + "payload", + BasicType.STRING_TYPE, + (Long) null, + true, + null, + null)) + .column( + PhysicalColumn.of( + "KafkaOffset", + BasicType.LONG_TYPE, + (Long) null, + true, + null, + "physical lookalike")) + .build(), + new HashMap<>(), + new ArrayList<>(), + "comment", + "test", + MetadataSchema.builder().build()); + + TransformException exception = + Assertions.assertThrows( + TransformException.class, + () -> + new MetadataTransform( + ReadonlyConfig.fromMap( + metadataFieldsConfig( + "KafkaOffset", "kafka_offset")), + table)); + Assertions.assertTrue(exception.getMessage().contains("KafkaOffset")); + } + + @Test + void shouldProjectNullConnectorMetadataValue() { + CatalogTable table = + catalogTableWithCustomMetadata( + Collections.singletonList( + MetadataColumn.of( + "KafkaOffset", + BasicType.LONG_TYPE, + (Long) null, + true, + null, + "Kafka record offset"))); + MetadataTransform transform = + new MetadataTransform( + ReadonlyConfig.fromMap(metadataFieldsConfig("KafkaOffset", "kafka_offset")), + table); + transform.initRowContainerGenerator(); + + SeaTunnelRow input = new SeaTunnelRow(new Object[] {"payload"}); + input.getOptions().put("KafkaOffset", null); + + SeaTunnelRow output = transform.map(input); + Assertions.assertEquals(2, output.getArity()); + Assertions.assertEquals("payload", output.getField(0)); + Assertions.assertNull(output.getField(1)); + } + + @Test + void shouldRejectConnectorMetadataKeyCaseMismatch() { + CatalogTable table = + catalogTableWithCustomMetadata( + Collections.singletonList( + MetadataColumn.of( + "KafkaOffset", + BasicType.LONG_TYPE, + (Long) null, + true, + null, + "Kafka record offset"))); + + TransformException exception = + Assertions.assertThrows( + TransformException.class, + () -> + new MetadataTransform( + ReadonlyConfig.fromMap( + metadataFieldsConfig( + "kafkaoffset", "kafka_offset")), + table)); + Assertions.assertTrue(exception.getMessage().contains("kafkaoffset")); + } + + @Test + void shouldProjectMultipleConnectorDeclaredFields() { + List<Column> metadata = new ArrayList<>(); + metadata.add( + MetadataColumn.of( + "KafkaOffset", + BasicType.LONG_TYPE, + (Long) null, + true, + null, + "Kafka record offset")); + metadata.add( + MetadataColumn.of( + "KafkaPartition", + BasicType.INT_TYPE, + (Long) null, + false, + 0, + "Kafka partition id")); + CatalogTable table = catalogTableWithCustomMetadata(metadata); + + Map<String, String> metadataMapping = new LinkedHashMap<>(); + metadataMapping.put("KafkaOffset", "kafka_offset"); + metadataMapping.put("KafkaPartition", "kafka_partition"); + Map<String, Object> config = new HashMap<>(); + config.put("metadata_fields", metadataMapping); + + MetadataTransform transform = new MetadataTransform(ReadonlyConfig.fromMap(config), table); + transform.initRowContainerGenerator(); + + Column[] columns = transform.getOutputColumns(); + Assertions.assertEquals("kafka_offset", columns[0].getName()); + Assertions.assertEquals(BasicType.LONG_TYPE, columns[0].getDataType()); + Assertions.assertTrue(columns[0].isNullable()); + Assertions.assertEquals("Kafka record offset", columns[0].getComment()); + Assertions.assertEquals("kafka_partition", columns[1].getName()); + Assertions.assertEquals(BasicType.INT_TYPE, columns[1].getDataType()); + Assertions.assertFalse(columns[1].isNullable()); + Assertions.assertEquals(0, columns[1].getDefaultValue()); + Assertions.assertEquals("Kafka partition id", columns[1].getComment()); + + SeaTunnelRow input = new SeaTunnelRow(new Object[] {"payload"}); + input.getOptions().put("KafkaOffset", 42L); + input.getOptions().put("KafkaPartition", 3); + + SeaTunnelRow output = transform.map(input); + Assertions.assertEquals(3, output.getArity()); + Assertions.assertEquals("payload", output.getField(0)); + Assertions.assertEquals(42L, output.getField(1)); + Assertions.assertEquals(3, output.getField(2)); + } + + @Test + void shouldKeepComputedAndCommonMetadataKeysUnchanged() { + Map<String, String> metadataMapping = new LinkedHashMap<>(); + metadataMapping.put("Database", "database"); + metadataMapping.put("Table", "table"); + metadataMapping.put("RowKind", "rowKind"); + metadataMapping.put("EventTime", "ts_ms"); + Map<String, Object> config = new HashMap<>(); + config.put("metadata_fields", metadataMapping); + + MetadataTransform transform = + new MetadataTransform(ReadonlyConfig.fromMap(config), catalogTable); + transform.initRowContainerGenerator(); + + Column[] columns = transform.getOutputColumns(); + Assertions.assertEquals(BasicType.STRING_TYPE, columns[0].getDataType()); + Assertions.assertEquals(BasicType.STRING_TYPE, columns[1].getDataType()); + Assertions.assertEquals(BasicType.STRING_TYPE, columns[2].getDataType()); + Assertions.assertEquals(BasicType.LONG_TYPE, columns[3].getDataType()); + + SeaTunnelRow outputRow = transform.map(inputRow); + Assertions.assertEquals("default", outputRow.getField(5)); + Assertions.assertEquals("default", outputRow.getField(6)); + Assertions.assertEquals("+I", outputRow.getField(7)); + Assertions.assertEquals(eventTime, outputRow.getField(8)); + } + @Test void shouldRejectMarkdownCanonicalPhysicalNameCollision() { Map<String, String> metadataMapping = new LinkedHashMap<>(); @@ -300,6 +546,34 @@ public class MetadataTransformTest { exception.getMessage().contains(KnowledgeSyncMetadataField.DOCUMENT_ID.getName())); } + private static Map<String, Object> metadataFieldsConfig(String metadataKey, String outputName) { + Map<String, String> metadataMapping = new LinkedHashMap<>(); + metadataMapping.put(metadataKey, outputName); + Map<String, Object> config = new HashMap<>(); + config.put("metadata_fields", metadataMapping); + return config; + } + + private static CatalogTable catalogTableWithCustomMetadata(List<Column> metadataColumns) { + return CatalogTable.of( + TableIdentifier.of("catalog", TablePath.DEFAULT), + TableSchema.builder() + .column( + PhysicalColumn.of( + "payload", + BasicType.STRING_TYPE, + (Long) null, + true, + null, + null)) + .build(), + new HashMap<>(), + new ArrayList<>(), + "comment", + "test", + MetadataSchema.builder().columns(metadataColumns).build()); + } + private static CatalogTable knowledgeSyncCatalogTable(boolean includeKnowledgeSyncMetadata) { List<Column> metadata = new ArrayList<>(); if (includeKnowledgeSyncMetadata) {
