{"record":{"id":"2461737de242553c","repo":"apache/druid","slug":"attempt-to-add-row-to-swapped-out-sink-for-segment","errorCode":null,"errorMessage":"Attempt to add row to swapped-out sink for segment[%s].","messagePattern":"Attempt to add row to swapped-out sink for segment\\[(.+?)\\]\\.","errorType":"exception","errorClass":"SegmentNotWritableException","httpStatus":null,"severity":"error","filePath":"server/src/main/java/org/apache/druid/segment/realtime/appenderator/StreamAppenderator.java","lineNumber":338,"sourceCode":"    }\n\n    final Sink sink = getOrCreateSink(identifier);\n    metrics.reportMessageMaxTimestamp(row.getTimestampFromEpoch());\n    final int sinkRowsInMemoryBeforeAdd = sink.getNumRowsInMemory();\n    final int sinkRowsInMemoryAfterAdd;\n    final long bytesInMemoryBeforeAdd = sink.getBytesInMemory();\n    final long bytesInMemoryAfterAdd;\n    final IncrementalIndexAddResult addResult;\n\n    addResult = sink.add(row);\n    sinkRowsInMemoryAfterAdd = addResult.getRowCount();\n    bytesInMemoryAfterAdd = addResult.getBytesInMemory();\n\n    final long currTs = System.currentTimeMillis();\n    metrics.reportMessageGap(currTs - row.getTimestampFromEpoch());\n\n    if (sinkRowsInMemoryAfterAdd < 0) {\n      throw new SegmentNotWritableException(\"Attempt to add row to swapped-out sink for segment[%s].\", identifier);\n    }\n\n    if (addResult.isRowAdded()) {\n      rowIngestionMeters.incrementProcessed();\n    } else if (addResult.hasParseException()) {\n      parseExceptionHandler.handle(addResult.getParseException());\n    }\n\n    final int numAddedRows = sinkRowsInMemoryAfterAdd - sinkRowsInMemoryBeforeAdd;\n    rowsCurrentlyInMemory.addAndGet(numAddedRows);\n    bytesCurrentlyInMemory.addAndGet(bytesInMemoryAfterAdd - bytesInMemoryBeforeAdd);\n    totalRows.addAndGet(numAddedRows);\n\n    boolean isPersistRequired = false;\n    boolean persist = false;\n    List<String> persistReasons = new ArrayList<>();\n\n    if (!sink.canAppendRow()) {","sourceCodeStart":320,"sourceCodeEnd":356,"githubUrl":"https://github.com/apache/druid/blob/9b90983fd291f26935af934383ce360473179e4d/server/src/main/java/org/apache/druid/segment/realtime/appenderator/StreamAppenderator.java#L320-L356","documentation":"StreamAppenderator.add throws SegmentNotWritableException when the target sink has been swapped out (persisted and replaced by an immutable hydrant), making it impossible to add more rows in memory to that segment. sinkRowsInMemoryAfterAdd < 0 is the sentinel indicating the sink's writable portion is gone.","triggerScenarios":"Adding a row to a segment whose sink was already swapped out mid-persist — a race between background persist/metadata publish and row arrival, or adding to a segment identifier after its sink was handed off.","commonSituations":"Kafka/Kinesis tasks during segment handoff where late-arriving records target a segment that was just persisted; replaying old offsets after a publish completed.","solutions":["Re-route the row to a newly allocated segment identifier (allocate a new pending segment for the event timestamp).","Check for clock/timestamp skew causing events to target old, already-persisted segments.","Upgrade/verify indexer versions — races between persist and add were fixed in later releases."],"exampleFix":"// before\nappenderator.add(oldIdentifier, row, supplier, false); // sink already swapped\n\n// after\nSegmentIdWithShardSpec id = allocateNewPendingSegment(row.getTimestamp());\nappenderator.add(id, row, supplier, false);","handlingStrategy":"try-catch","validationCode":"// skip add when the sink was handed off\nif (appenderator.getSinks().get(identifier) == null || !isSinkWritable(identifier)) { allocateNewSegmentAndRetry(row); return; }","typeGuard":null,"tryCatchPattern":"try { appenderator.add(id, row, supplier, true) } catch (SegmentNotWritableException e) { SegmentIdWithShardSpec fresh = allocator.newSegment(row.getTimestamp()); appenderator.add(fresh, row, supplier, true); }","preventionTips":["Keep event timestamps within the active segment window","Handle late events by allocating new pending segments","Avoid replaying already-published offsets"],"tags":["druid","ingestion","segment","concurrency"],"backgroundTag":"invalid-state-transition","analyzedSha":"9b90983fd291f26935af934383ce360473179e4d","analyzedAt":"2026-09-07T13:32:30.957Z","contentChangedAt":"2026-09-07T13:32:30.957Z","schemaVersion":2},"datasetVersion":"2026-09-14T05:17:10.506Z"}