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
