Showing posts with label nifi. Show all posts
Showing posts with label nifi. Show all posts

Thursday, 24 September 2026

Understanding NiFi Processors: The Engines Behind Every Data Flow

  

When you first start working with Apache NiFi, one of the most overwhelming experiences is opening the Processor dialog and seeing hundreds of available processors. New users often wonder:

 

·      Which processor should I use?

·      What is the difference between GetFile and FetchFile?

·      How do I transform data?

·      How do I route FlowFiles?

·      How do I send data to Kafka, databases, or cloud services?

 

The good news is that you do not need to learn every processor individually. Most processors fall into a few common categories such as data ingestion, transformation, routing, enrichment, database access, system integration, and data delivery.

 

In this post, we will explore the major processor categories, understand their purpose, and review commonly used processors that appear in real-world NiFi data flows.

 


1. What is a Processor?

A Processor is the fundamental building block of an Apache NiFi data flow. Think of a Processor as a worker in an assembly line:

 

·      Receives a FlowFile

·      Performs an action

·      Sends the FlowFile to the next processor

 

Examples:

·      Read files

·      Execute SQL queries

·      Convert JSON

·      Route data

·      Send messages to Kafka

·      Upload files to S3

 

A complete NiFi flow is simply multiple processors connected together.

GetFile
   |
   v
SplitJson
   |
   v
RouteOnAttribute
   |
   +------> PublishKafka
   |
   +------> PutFile

   

2. Processor Categories

Although NiFi contains hundreds of processors, most belong to the following categories.

 

2.1 Data Ingestion Processors

These processors bring data into NiFi.

 

Processor

Purpose

GetFile     

Read files from local disk    

GetFTP      

Download files from FTP       

GetSFTP     

Download files from SFTP      

ConsumeKafka

Read messages from Kafka      

ConsumeJMS  

Read messages from JMS        

ListenHTTP  

Receive HTTP requests         

ListenTCP   

Receive TCP messages          

ListenUDP   

Receive UDP packets           

TailFile    

Continuously monitor log files

GetMongo    

Read data from MongoDB        

GetSQS      

Read messages from AWS SQS    

 

2.2 Data Transformation Processors

These processors change the content or format of data.

 

Processor

Purpose

ConvertRecord      

Convert CSV ↔ JSON ↔ Avro

JoltTransformJSON  

Transform JSON               

JSLTTransformJSON  

Advanced JSON transformations

TransformXml       

XML transformation using XSLT

ReplaceText        

Replace text using regex     

FlattenJson        

Flatten nested JSON          

ConvertCharacterSet

Convert encoding formats     

 

2.3 Routing Processors

These processors decide where data should go.

Processor

Purpose

RouteOnAttribute

Route based on attributes

RouteOnContent  

Route based on content   

RouteText       

Route text patterns      

ScanAttribute   

Search attributes        

ScanContent     

Search content           

ValidateJson    

Validate JSON            

ValidateXml     

Validate XML             

 

2.4 Attribute Management Processors

FlowFile attributes store metadata.

 

Processor

Purpose

UpdateAttribute         

Add or modify attributes           

EvaluateJsonPath        

Extract JSON values into attributes

EvaluateXPath           

Extract XML values                 

ExtractText             

Extract text using regex           

LookupAttribute         

Lookup attribute values            

CryptographicHashContent

Generate content hash              

IdentifyMimeType        

Detect file type                   

 

2.5 Database Processors

These processors interact with relational databases.

 

Processor

Purpose

ExecuteSQL        

Execute SELECT query           

ExecuteSQLRecord  

Execute query using Record API 

PutSQL            

Execute INSERT/UPDATE/DELETE   

PutDatabaseRecord 

Write records into database    

QueryDatabaseTable

Incremental database polling   

GenerateTableFetch

Generate parallel fetch queries

 

2.6 Splitting and Merging Processors

Large files often need to be split before processing.

 

Processor

Purpose

SplitJson     

Split JSON arrays      

SplitXml      

Split XML              

SplitText     

Split text files       

SplitRecord   

Split record-based data

SegmentContent

Split by size          

SplitContent  

Split using delimiter  

 

Merging Processors

Processor

Purpose

MergeContent

Merge files         

MergeRecord 

Merge record streams

 

2.7 Compression and Encryption Processors

Used for security and storage optimization.

 

Processor

Purpose

CompressContent  

Compress data     

UnpackContent    

Extract ZIP/TAR   

EncryptContentPGP

Encrypt content   

DecryptContentPGP

Decrypt content   

SignContentPGP   

Digital signatures

VerifyContentPGP

Verify signatures 

 

2.8 System Integration Processors

These processors interact with operating systems and external scripts.

 

Processor

Purpose

ExecuteProcess         

Execute OS command     

ExecuteStreamCommand   

Stream data to commands

ExecuteScript          

Execute scripts        

ExecuteGroovyScript    

Run Groovy code        

InvokeScriptedProcessor

Custom scripting       

 

2.9 Messaging and Streaming Processors

These processors connect NiFi with messaging systems.

 

Processor

Purpose

ConsumeKafka

Consumes messages from Apache Kafka Consumer API.

PublishKafka

Sends the contents of a FlowFile as either a message or as individual records to Apache Kafka using the Kafka Producer API.

ConsumeJMS

Consumes JMS Message of type BytesMessage, TextMessage, ObjectMessage, MapMessage or StreamMessage transforming its content to a FlowFile and transitioning it to 'success' relationship.

PublishJMS

Creates a JMS Message from the contents of a FlowFile and sends it to a JMS Destination (queue or topic) as JMS BytesMessage or TextMessage. FlowFile attributes will be added as JMS headers and/or properties to the outgoing JMS message.

ConsumeMQTT

Subscribes to a topic and receives messages from an MQTT broker

PublishMQTT

Publishes a message to an MQTT topic

ConsumeAMQP

Consumes AMQP Messages from an AMQP Broker using the AMQP 0.9.1 protocol. Each message that is received from the AMQP Broker will be emitted as its own FlowFile to the 'success' relationship.

PublishAMQP

Creates an AMQP Message from the contents of a FlowFile and sends the message to an AMQP Exchange.

 

 

2.10 Cloud Processors

 

NiFi provides processors for major cloud providers.

 

Processor

Purpose

FetchS3Object

Retrieves the contents of an S3 Object and writes it to the content of a FlowFile

PutS3Object

Writes the contents of a FlowFile as an S3 Object to an Amazon S3 Bucket.

GetSQS

Fetches messages from an Amazon Simple Queuing Service Queue

PutSQS

Publishes a message to an Amazon Simple Queuing Service Queue

PutSNS

Sends the content of a FlowFile as a notification to the Amazon Simple Notification Service

PutLambda

Sends the contents to a specified Amazon Lambda Function.

PutKinesisStream

Sends the contents to a specified Amazon Kinesis.

GetAzureEventHub

Receives messages from Microsoft Azure Event Hubs without reliable checkpoint tracking.

PutAzureEventHub

Send FlowFile contents to Azure Event Hubs

PutAzureBlobStorage

Puts content into a blob on Azure Blob Storage.

PutAzureDataLakeStorage

Writes the contents of a FlowFile as a file on Azure Data Lake Storage Gen 2

ConsumeGCPubSub

Consumes messages from the configured Google Cloud PubSub subscription.

PublishGCPubSub

Publishes the content of the incoming flowfile to the configured Google Cloud PubSub topic.

PutBigQuery

Writes the contents of a FlowFile to a Google BigQuery table.

PutGCSObject

Writes the contents of a FlowFile as an object in a Google Cloud Storage.

 

2.11 Monitoring and Control Processors

 

These processors help manage and monitor data flows.

 

Processor

Purpose

MonitorActivity

Detect inactivity      

ControlRate    

Throttle flow          

DetectDuplicate

Duplicate detection    

RetryFlowFile  

Retry failed processing

Notify         

Trigger notifications  

Wait           

Synchronization        

DebugFlow      

Debug data flow        

 

2.12 Record-Oriented Processors

Modern NiFi development heavily uses Record APIs because they provide better performance and flexibility.

 

Processor

Purpose

ConvertRecord    

Converts records from one data format to another using configured Record Reader and Record Write Controller Services.     

QueryRecord      

Evaluates one or more SQL queries against the contents of a FlowFile.

LookupRecord     

Extracts one or more fields from a Record and looks up a value for those fields in a LookupService.

UpdateRecord     

Updates the contents of a FlowFile that contains Record-oriented data

ValidateRecord   

Validates the Records of an incoming FlowFile against a given schema.

MergeRecord      

This Processor merges together multiple record-oriented FlowFiles into a single FlowFile that contains all of the Records of the input FlowFiles.

PartitionRecord  

Splits, or partitions, record-oriented data based on the configured fields in the data.

DeduplicateRecord

This processor de-duplicates individual records within a record set.

 

A Simple Real-World Example

Imagine a company receives customer orders through SFTP and wants to publish them to Kafka.

 

GetSFTP
    |
    v
SplitJson
    |
    v
ValidateJson
    |
    v
UpdateAttribute
    |
    v
PublishKafka

   

Processor roles:

 

·      GetSFTP: Fetch orders

·      SplitJson: Separate individual orders

·      ValidateJson: Validate structure

·      UpdateAttribute: Add metadata

·      PublishKafka: Send downstream

 

In summary, NiFi contains hundreds of processors, but most fall into a small number of categories. Once you understand these categories, navigating NiFi's extensive processor library becomes much easier, and you can quickly identify the right processor for almost any data integration requirement.

 

Previous                                                    Next                                                    Home