You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何配置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) or LocalMapCacheClientService (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)

Step 3: Check the Cache with FetchDistributedMapCache

  • Link the success relationship from UpdateAttribute to 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'"
  • 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 miss relationship from FetchDistributedMapCache to your target processor (the one you want to run once daily).
  • Add a PutDistributedMapCache processor, and connect your target processor's success relationship 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.
  • Route the hit relationship from FetchDistributedMapCache to a TerminateFlowFile processor—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 ExecuteGroovyScript processor 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 success to your target processor, and failure to TerminateFlowFile.

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.06 13:59:05