{"record":{"id":"fa8daa18aeba4058","repo":"apache/iceberg","slug":"failed-to-serialize-pk-index-key-fa8daa","errorCode":null,"errorMessage":"Failed to serialize PK index key","messagePattern":"Failed to serialize PK index key","errorType":"exception","errorClass":"UncheckedIOException","httpStatus":null,"severity":"error","filePath":"flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/StructLikeSerializer.java","lineNumber":68,"sourceCode":"  private final ByteArrayOutputStream baos = new ByteArrayOutputStream();\n  private final DataOutputStream dos = new DataOutputStream(baos);\n\n  public SerializedEqualityValues serializeKey(StructLike key, Types.StructType keyType) {\n    baos.reset();\n    try {\n      List<Types.NestedField> fields = keyType.fields();\n      dos.writeInt(fields.size());\n      for (Types.NestedField field : fields) {\n        dos.writeInt(field.fieldId());\n      }\n\n      for (int i = 0; i < fields.size(); i++) {\n        writeField(key, i, fields.get(i).type());\n      }\n\n      dos.flush();\n    } catch (IOException e) {\n      throw new UncheckedIOException(\"Failed to serialize PK index key\", e);\n    }\n\n    return new SerializedEqualityValues(baos.toByteArray());\n  }\n\n  public byte[] encodePartition(StructLike partition, Types.StructType partitionType) {\n    List<Types.NestedField> fields = partitionType.fields();\n    if (fields.isEmpty()) {\n      return EMPTY_PARTITION;\n    }\n\n    baos.reset();\n    try {\n      for (int i = 0; i < fields.size(); i++) {\n        writeField(partition, i, fields.get(i).type());\n      }\n\n      dos.flush();","sourceCodeStart":50,"sourceCodeEnd":86,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/StructLikeSerializer.java#L50-L86","documentation":"serializeKey() writes the equality-field ids and values of a primary-key (equality) key into Flink keyed state. Any IOException while writing the field bytes is wrapped in UncheckedIOException. In practice this means the underlying StructLike could not be converted for one of the declared key field types (e.g. an unexpected Java type for the field).","triggerScenarios":"Calling serializeKey with a StructLike whose field values do not match keyType's declared types so Conversions.toByteBuffer throws; a corrupted struct whose get() returns an unconvertible object.","commonSituations":"Row schema evolved between jobs so the equality key type no longer matches the stored rows; custom StructLike implementations returning unexpected types; upstream produced rows that violate the declared schema.","solutions":["Verify the equality key field types match the actual values in the DataStream/schema; fix the schema definition passed to the maintenance source","Check for schema evolution (added/renamed key fields) between the writing job and the current job and rebuild state","Log the failing key's field values before this call in a debug run to find which field/value is incompatible","Update to a newer Iceberg version if a type conversion gap for your field type was fixed"],"exampleFix":"// before: key field declared as TimestampType but row supplies LocalDateTime\nRowType keyType = RowType.of(new TimestampType(), ...);\n// after: declare the key type matching the supplied Java value\nRowType keyType = RowType.of(new LocalZonedTimestampType(), ...);","handlingStrategy":"validation","validationCode":"for (int i = 0; i < keyType.fields().size(); i++) {\n  Object v = key.get(i, Object.class);\n  Preconditions.checkState(v == null || Conversions.tryToByteBuffer(keyType.fields().get(i).type(), v) != null,\n      \"Key field %s value %s not convertible\", keyType.fields().get(i).name(), v);\n}","typeGuard":null,"tryCatchPattern":"try {\n  serializer.serializeKey(key, keyType);\n} catch (UncheckedIOException e) {\n  LOG.error(\"PK key serialization failed for key {}\", key, e);\n  throw e;\n}","preventionTips":["Keep equality key field types in sync with the RowType/schema used by the job","Rebuild keyed state after schema/key evolution instead of restoring old state","Test serialization round-trips of keys in unit tests when changing schemas"],"tags":["flink","serialization","unchecked-io","structlike"],"backgroundTag":"json-serialization-failed","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"}