{"record":{"id":"66803fb7c7028894","repo":"apache/iceberg","slug":"record-value-is-not-valid-json-for-record-value","errorCode":null,"errorMessage":"record.value is not valid json for record.value: ${collectRecordDetails(record)}","messagePattern":"record\\.value is not valid json for record\\.value: (.+?)","errorType":"exception","errorClass":"JsonToMapException","httpStatus":null,"severity":"error","filePath":"kafka-connect/kafka-connect-transforms/src/main/java/org/apache/iceberg/connect/transforms/JsonToMapTransform.java","lineNumber":82,"sourceCode":"    if (record.value() == null) {\n      return record;\n    } else {\n      return process(record);\n    }\n  }\n\n  private R process(R record) {\n    if (!(record.value() instanceof String)) {\n      throw new JsonToMapException(\"record value is not a string, use StringConverter\");\n    }\n\n    String json = (String) record.value();\n    JsonNode obj;\n\n    try {\n      obj = MAPPER.readTree(json);\n    } catch (Exception e) {\n      throw new JsonToMapException(\n          String.format(\n              \"record.value is not valid json for record.value: %s\", collectRecordDetails(record)),\n          e);\n    }\n\n    if (!(obj instanceof ObjectNode)) {\n      throw new JsonToMapException(\n          String.format(\n              \"Expected json object for record.value after parsing: %s\",\n              collectRecordDetails(record)));\n    }\n\n    if (startAtRoot) {\n      return singleField(record, (ObjectNode) obj);\n    }\n    return structRecord(record, (ObjectNode) obj);\n  }\n","sourceCodeStart":64,"sourceCodeEnd":100,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/kafka-connect/kafka-connect-transforms/src/main/java/org/apache/iceberg/connect/transforms/JsonToMapTransform.java#L64-L100","documentation":"After confirming the value is a String, JsonToMapTransform parses it with Jackson readTree. If parsing fails (malformed JSON), it wraps the cause in JsonToMapException including record details ('record.value is not valid json for record.value: ...'). This means the value string is not syntactically valid JSON.","triggerScenarios":"record.value() is a String but contains truncated output, non-JSON text (log lines, CSV), encoding corruption, or double-encoded/malformed JSON from the producer.","commonSituations":"Producers writing plain-text or CSV to a topic consumed by Iceberg sink; message size truncation; DMS emitting non-JSON payloads; charset/encoding mismatches corrupting bytes.","solutions":["Fix the producer to emit valid JSON strings","Validate sample messages from the topic (kafka-console-consumer) before configuring the sink","Check for truncation (max.message.bytes / connector buffer limits) and encoding issues","Route malformed messages to a dead-letter queue via errors.tolerance=all + errors.deadletterqueue.*"],"exampleFix":"// before (connector config, fail on bad data)\n\"errors.tolerance\": \"none\"\n// after (quarantine bad records for inspection)\n\"errors.tolerance\": \"all\",\n\"errors.deadletterqueue.topic.name\": \"dlq-iceberg\",\n\"errors.deadletterqueue.context.headers.enable\": \"true\"","handlingStrategy":"validation","validationCode":"boolean isValidJson(String s) {\n  try { new ObjectMapper().readTree(s); return true; } catch (Exception e) { return false; }\n}","typeGuard":null,"tryCatchPattern":"try {\n  return transform.apply(record);\n} catch (JsonToMapException e) {\n  deadLetterQueue.send(record, e); // with errors.tolerance=all\n  return null;\n}","preventionTips":["Fix producers to emit valid JSON","Enable DLQ: errors.tolerance=all + errors.deadletterqueue.topic.name","Sample topic data before configuring the Iceberg sink","Check for truncation and encoding corruption in the transport"],"tags":["kafka-connect","smt","json","parse-error"],"backgroundTag":"json-parse-error","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"}