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();

Reply via email to