{"record":{"id":"32123b4a88cf64bd","repo":"apache/iceberg","slug":"schema-not-support-for-dms-records","errorCode":null,"errorMessage":"Schema not support for DMS records","messagePattern":"Schema not support for DMS records","errorType":"exception","errorClass":"UnsupportedOperationException","httpStatus":null,"severity":"error","filePath":"kafka-connect/kafka-connect-transforms/src/main/java/org/apache/iceberg/connect/transforms/DmsTransform.java","lineNumber":42,"sourceCode":"import org.apache.kafka.connect.connector.ConnectRecord;\nimport org.apache.kafka.connect.transforms.Transformation;\nimport org.apache.kafka.connect.transforms.util.Requirements;\nimport org.slf4j.Logger;\nimport org.slf4j.LoggerFactory;\n\npublic class DmsTransform<R extends ConnectRecord<R>> implements Transformation<R> {\n\n  private static final Logger LOG = LoggerFactory.getLogger(DmsTransform.class.getName());\n  private static final ConfigDef EMPTY_CONFIG = new ConfigDef();\n\n  @Override\n  public R apply(R record) {\n    if (record.value() == null) {\n      return record;\n    } else if (record.valueSchema() == null) {\n      return applySchemaless(record);\n    } else {\n      throw new UnsupportedOperationException(\"Schema not support for DMS records\");\n    }\n  }\n\n  @SuppressWarnings(\"unchecked\")\n  private R applySchemaless(R record) {\n    Map<String, Object> value = Requirements.requireMap(record.value(), \"DMS transform\");\n\n    // promote fields under \"data\"\n    Object dataObj = value.get(\"data\");\n    Object metadataObj = value.get(\"metadata\");\n    if (!(dataObj instanceof Map) || !(metadataObj instanceof Map)) {\n      LOG.debug(\"Unable to transform DMS record, skipping...\");\n      return null;\n    }\n\n    Map<String, Object> metadata = (Map<String, Object>) metadataObj;\n\n    String dmsOp = metadata.get(\"operation\").toString();","sourceCodeStart":24,"sourceCodeEnd":60,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/kafka-connect/kafka-connect-transforms/src/main/java/org/apache/iceberg/connect/transforms/DmsTransform.java#L24-L60","documentation":"The DmsTransform SMT only supports schemaless Kafka Connect records: if record.value() is non-null and record.valueSchema() is non-null (a record WITH schema), it throws UnsupportedOperationException('Schema not support for DMS records'). AWS DMS records are expected as raw JSON map values without a Connect schema.","triggerScenarios":"Configuring value.converter=JsonConverter with schemas.enable=true (or AvroConverter) upstream of DmsTransform, so records arrive with a non-null valueSchema.","commonSituations":"AWS DMS → Kafka Connect pipelines where the value converter emits schemas; switching converters from schemaless JSON to schema-bearing formats; Kafka Connect defaults where schemas.enable=true is set globally.","solutions":["Set value.converter.schemas.enable=false on the connector so values are schemaless maps","Switch the value converter to org.apache.kafka.connect.json.JsonConverter with schemas.enable=false","Remove DmsTransform and use a different transform if you need schema-bearing records","Preprocess records to strip value schemas before the transform"],"exampleFix":"// before (worker/connector config)\n\"value.converter.schemas.enable\": \"true\"\n// after\n\"value.converter\": \"org.apache.kafka.connect.json.JsonConverter\",\n\"value.converter.schemas.enable\": \"false\"","handlingStrategy":"validation","validationCode":"boolean dmsTransformCompatible(SinkRecord r) {\n  return r.value() == null || r.valueSchema() == null;\n}","typeGuard":"boolean isSchemaless(SinkRecord r) {\n  return r.valueSchema() == null && r.value() instanceof Map;\n}","tryCatchPattern":"try {\n  return dmsTransform.apply(record);\n} catch (UnsupportedOperationException e) {\n  throw new ConnectException(\"DmsTransform requires schemaless values; set value.converter.schemas.enable=false\", e);\n}","preventionTips":["Set value.converter.schemas.enable=false when using DmsTransform","Avoid AvroConverter/ProtobufConverter upstream of this SMT","Unit-test the SMT chain with the actual converter config"],"tags":["kafka-connect","smt","dms","config"],"backgroundTag":"unsupported-operation","analyzedSha":"86d9c8fc543e7c56c9f624eb725f76c9baff9570","analyzedAt":"2026-09-12T00:46:39.097Z","contentChangedAt":"2026-09-12T00:46:39.097Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}