This is an automated email from the ASF dual-hosted git repository.
pvillard31 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/nifi.git
The following commit(s) were added to refs/heads/main by this push:
new 7b2580042c6 NIFI-15904 Fixed Lineage Start Index in Session.create()
(#11205)
7b2580042c6 is described below
commit 7b2580042c616ecbe4d52e2924aba70f1cb45e84
Author: David Handermann <[email protected]>
AuthorDate: Tue May 5 03:47:28 2026 -0500
NIFI-15904 Fixed Lineage Start Index in Session.create() (#11205)
- Set explicit Lineage Start Index using FlowFile ID
- Set explicit Entry Date and Lineage Start Date on new FlowFiles
---
.../repository/StandardProcessSession.java | 7 ++++++-
.../repository/StandardProcessSessionTest.java | 22 ++++++++++++++++++++++
2 files changed, 28 insertions(+), 1 deletion(-)
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/StandardProcessSession.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/StandardProcessSession.java
index a97e47326cc..4ae94c47ca6 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/StandardProcessSession.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/StandardProcessSession.java
@@ -2069,8 +2069,13 @@ public class StandardProcessSession implements
ProcessSession, ProvenanceEventEn
attrs.put(CoreAttributes.PATH.key(), DEFAULT_FLOWFILE_PATH);
attrs.put(CoreAttributes.UUID.key(), uuid);
- final FlowFileRecord fFile = new
StandardFlowFileRecord.Builder().id(context.getNextFlowFileSequence())
+ final long entryDate = System.currentTimeMillis();
+ final long id = context.getNextFlowFileSequence();
+ final FlowFileRecord fFile = new
StandardFlowFileRecord.Builder().id(id)
.addAttributes(attrs)
+ // Set Lineage Start to Entry Date and use Identifier as Lineage
Start Index for unambiguous ordering
+ .entryDate(entryDate)
+ .lineageStart(entryDate, id)
.build();
final StandardRepositoryRecord record = new
StandardRepositoryRecord((FlowFileQueue) null);
record.setWorking(fFile, attrs, false);
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/controller/repository/StandardProcessSessionTest.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/controller/repository/StandardProcessSessionTest.java
index c82ee1ca577..42c2192a5ac 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/controller/repository/StandardProcessSessionTest.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/controller/repository/StandardProcessSessionTest.java
@@ -45,6 +45,8 @@ import java.nio.file.Files;
import java.nio.file.Path;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.anyString;
@@ -191,6 +193,26 @@ class StandardProcessSessionTest {
assertEquals(GAUGE_VALUE, gaugeRecord.value());
}
+ @Test
+ void testCreateLineage() {
+ final long firstFlowFileId = 1;
+ final long secondFlowFileId = 2;
+
when(repositoryContext.getNextFlowFileSequence()).thenReturn(firstFlowFileId,
secondFlowFileId);
+
+ final FlowFile firstFlowFile = session.create();
+
+ assertNotNull(firstFlowFile);
+ assertNotEquals(0, firstFlowFile.getLineageStartDate());
+ assertEquals(firstFlowFile.getEntryDate(),
firstFlowFile.getLineageStartDate());
+ assertEquals(firstFlowFileId, firstFlowFile.getId());
+ assertEquals(firstFlowFileId, firstFlowFile.getLineageStartIndex());
+
+ final FlowFile secondFlowFile = session.create();
+ assertNotNull(secondFlowFile);
+ assertEquals(secondFlowFileId, secondFlowFile.getId());
+ assertEquals(secondFlowFileId, secondFlowFile.getLineageStartIndex());
+ }
+
private void assertFlowFileEventMatched(final long bytesRead, final long
bytesWritten) throws IOException {
verify(flowFileEventRepository).updateRepository(flowFileEventCaptor.capture(),
anyString());
final FlowFileEvent flowFileEvent = flowFileEventCaptor.getValue();