如何配置NiFi处理器仅在接收当日首个FlowFile时运行并忽略后续文件?
Got it, let's figure out how to make your NiFi processor run only on the first FlowFile of the day and ignore everything that comes after. I've got two solid approaches for you—one using built-in components (no custom code needed) and another with a quick script for single-node setups:
Approach 1: Built-in NiFi Components (Clustered Friendly)
This method uses NiFi's cache and routing processors to track whether we've processed the first FlowFile of the day, perfect for multi-node clusters.
Step 1: Set Up a Cache Service
- Add either a
DistributedMapCacheClientService(for clusters) orLocalMapCacheClientService(single node) to your canvas. - Give it a unique cache name (like
DailyFirstFlowTracker) and enable the service.
Step 2: Add an UpdateAttribute Processor
- Connect your FlowFile source to this processor.
- Add a new attribute to capture today's date:
- Name:
current_day - Value:
${now():format('yyyy-MM-dd')}(this creates a consistent, sortable date string)
- Name:
Step 3: Check the Cache with FetchDistributedMapCache
- Link the
successrelationship fromUpdateAttributeto this processor. - Configure these key settings:
- Cache Entry Identifier:
${current_day} - Cache Service: Select the cache service you set up.
- If Cache Entry Not Found: Set to "Route to 'miss'"
- Cache Entry Identifier:
- This processor will check if today's date exists in the cache (meaning we've already processed a FlowFile today).
Step 4: Route FlowFiles and Update Cache
- Send the
missrelationship fromFetchDistributedMapCacheto your target processor (the one you want to run once daily). - Add a
PutDistributedMapCacheprocessor, and connect your target processor'ssuccessrelationship to it. - Configure
PutDistributedMapCache:- Cache Entry Identifier:
${current_day} - Cache Value:
processed(any non-empty value works) - Cache Service: Same as your earlier cache setup.
- Cache Entry Identifier:
- Route the
hitrelationship fromFetchDistributedMapCacheto aTerminateFlowFileprocessor—this will discard all subsequent FlowFiles for the day.
Approach 2: ExecuteGroovyScript (Single Node Simplified)
If you're running a single NiFi node, you can skip the cache service and use a script to manage local processor state:
- Add an
ExecuteGroovyScriptprocessor and connect your FlowFile source to it. - Paste this script into the "Script" field:
import org.apache.nifi.processor.state.Scope def flowFile = session.get() if (!flowFile) return // Get today's date in a standard format def today = new Date().format('yyyy-MM-dd') // Fetch existing state to check if we've processed today def currentState = context.getStateManager().getState(Scope.LOCAL) def processedToday = currentState.get("last_processed_day") == today if (!processedToday) { // Update state to mark today as processed def newState = [:] newState.put("last_processed_day", today) context.getStateManager().setState(newState, Scope.LOCAL) // Send to your target processor session.transfer(flowFile, REL_SUCCESS) } else { // Discard the FlowFile session.transfer(flowFile, REL_FAILURE) }
- Keep the default relationships: route
successto your target processor, andfailuretoTerminateFlowFile.
Quick Tips
- For clusters, Approach 1 is mandatory—the distributed cache ensures only one FlowFile gets processed across all nodes each day.
- To reset manually: For Approach 1, clear the cache via the cache service's UI; for Approach 2, right-click the script processor -> View State -> Clear State.
内容的提问来源于stack exchange,提问作者denimmonkey
相关产品推荐
相关产品推荐

