A practical design is to use Apache NiFi to ingest and route events, Apache Pulsar to transport and retain them, and Apache Flink to process them before writing results to a chosen sink. But do not assume this is a plug-and-play integration: the cited NiFi guide documents Kafka processors, not a verified native Pulsar processor, and Flink’s current Pulsar connector documentation says a SQL JAR is not available for Flink 2.3. Confirm a compatible NiFi-to-Pulsar path and Flink connector for your exact releases before building around this design.
What the architecture does—and what must be verified
These projects cover different parts of a streaming system. NiFi provides a processor-driven flow for moving and routing FlowFiles; Pulsar provides topic-based messaging and retention; Flink processes streams. Apache Pulsar describes its platform as horizontally scalable, but that is a capability of the platform, not a capacity or latency guarantee for a particular application. See the Apache Pulsar project overview.
| Layer | Role in the application | Decision to make |
|---|---|---|
| Apache NiFi | Ingest, validate, transform, and route incoming data. | Select and test a Pulsar publishing mechanism compatible with your NiFi release. The NiFi Getting Started guide documents Kafka processors; it does not establish an equivalent supported Pulsar processor. |
| Apache Pulsar | Carry events between producers and consumers, organized into topics. | Choose a release line, topic and partition plan, retention policy, schema/serialization, and security configuration. Use documentation matching that release, such as the Pulsar 5.0 documentation portal. |
| Apache Flink | Consume events, perform stateful or stateless stream processing, and emit results. | Verify the connector artifact and whether the desired Flink release supports direct SQL access to Pulsar. Connector dependencies are separate from the Flink binary distribution. |
| Output sink | Store or deliver processed results, for example to a database or another messaging system. | Choose a sink whose retry and idempotency behavior fits the delivery guarantees the application needs. |
This is a reference topology, not a prebuilt, currently verified NiFi → Pulsar → Flink SQL integration. Apache Pulsar’s docs are published by release line, and connector compatibility is version-specific. Pin NiFi, Pulsar, Flink, and connector versions together in your deployment plan instead of mixing instructions from different release lines.
How to connect Apache NiFi to Apache Pulsar
First settle the publishing path; the cited NiFi documentation alone does not establish a native Pulsar processor for a particular NiFi release. Verify the NiFi extension registry and the selected extension’s compatibility, security support, delivery behavior, and maintenance status. If you cannot verify an appropriate extension, use a documented and tested custom-client route or a separately operated bridge. Do not treat a Kafka processor as a Pulsar publisher without a verified compatibility layer.
Crashes, No Sound, or Screen Glitches?
Random freezes, missing sound and display glitches usually trace back to one bad driver. Find and replace yours safely.Free scan · under a minutePC Slower Than It Used to Be?
A free scan shows the junk files, broken settings and background clutter dragging Windows down - then fixes them in one click.Free scan · Windows 10 & 11#1 Best Overall
- Choose the exact NiFi and Pulsar releases. Use the matching Pulsar documentation and confirm that the chosen NiFi extension or client supports those versions.
- Define how NiFi publishes. Document the processor or client, Pulsar service endpoint, authentication and TLS settings, target topic, serialization, and error handling. Test the path with the actual deployment configuration.
- Define failure handling. Decide what NiFi should do when Pulsar is unavailable or rejects a record: retry, retain it in a failure route, or apply another explicit recovery policy. Test the chosen behavior rather than assuming delivery semantics.
- Verify with a consumer. Publish representative valid and malformed events, then confirm the expected topic, payload, key, and failure route from a consumer before connecting Flink.
Can Flink SQL read from Pulsar?
Do not assume direct Flink SQL-to-Pulsar support. The Flink Pulsar connector documentation describes a DataStream connector and states that there is no SQL JAR available for Flink 2.3. Because stable and nightly documentation can change, check the docs for the exact Flink release you intend to deploy before writing SQL DDL. The cited connector page also notes that connector dependencies are not bundled in the Flink binary distribution, so the matching dependency must be made available to the cluster.
If direct SQL support is not verified for your target release, there are two honest choices:
Rank #2
- Use a documented DataStream path. Consume Pulsar with a compatible Flink Pulsar connector and implement the processing in a Flink DataStream job. This is an alternative to Flink SQL, not SQL under another name.
- Change the version or integration design. Select a release and connector combination whose SQL support is explicitly documented, then validate it in a small deployment before committing to the architecture.
Do not silently substitute Kafka SQL, Pulsar SQL, or Trino and describe the result as Flink SQL. The Flink connector configuration described in its documentation includes the Pulsar service URL, admin URL, subscription name, topic or partition selection, and deserialization configuration. Supply the values and serialization appropriate to your cluster and data contract.
Design the event contract before processing
A pipeline can connect successfully yet still produce unusable or inconsistent results if producers and consumers disagree about event meaning. Define the contract at the topic boundary and make it part of deployment and change review.
Rank #3
- Topic layout: establish topic names and decide whether workloads or event types need separate topics. Choose partitioning and message keys based on ordering and parallelism needs.
- Schema and serialization: specify the payload format, required fields, compatibility policy for schema changes, and how consumers handle unknown or missing fields.
- Event time: distinguish the time an event occurred from the time it was ingested or processed. Decide how late or out-of-order events should affect results.
- Invalid data: define validation rules and a route or policy for malformed, unsupported, or incomplete records so failures are visible rather than silently discarded.
- Security: specify authentication, authorization, and transport security for NiFi-to-Pulsar and Flink-to-Pulsar connections; do not assume a connector inherits credentials or settings from another component.
How to make a Flink–Pulsar pipeline fault tolerant
Fault tolerance is a property of the whole path, not a single checkbox. Flink checkpointing, Pulsar subscription and acknowledgement behavior, connector configuration, and the output sink all affect recovery and duplicate handling.
Checkpointing and Pulsar consumption
Enable and test Flink checkpoints as part of the job’s recovery design. The Flink 2.1 Pulsar connector documentation describes source acknowledgement at completed checkpoints for documented subscription modes, and different acknowledgement behavior when checkpointing is disabled. It also describes transaction requirements for shared/key-shared use and warns that immediate acknowledgement without the relevant consistency setup does not provide the same guarantees. Select the subscription mode and connector configuration deliberately; do not label the pipeline “exactly once” without validating the complete configuration and failure cases.
Rank #4
Sink behavior and duplicates
After a failure, records can be retried or replayed. The sink must therefore be considered alongside the source: use an idempotent write strategy or a transactional sink where supported and required, and define how duplicate results are recognized or reconciled. Apache Pulsar’s connector overview notes that delivery guarantees depend on the connector and external sink implementation, including whether retries are idempotent.
Recovery tests
- Stop or make the Flink job unavailable during processing, then confirm it recovers from the configured checkpoint and resumes consumption.
- Interrupt Pulsar connectivity and confirm the job reports the failure and resumes without silently losing records under the chosen configuration.
- Force a sink failure and verify retries, duplicate behavior, and eventual output against the sink’s actual semantics.
- Test malformed events separately so they do not create an invisible failure loop or block unrelated valid traffic.
Scale each layer from measurements
There is no responsible universal broker count, partition count, NiFi concurrency, or Flink parallelism for this topology. Size it from the workload and then validate under representative conditions. Gather event rate and size, burst patterns, retention duration, latency objective, availability objective, schema characteristics, and deployment topology before choosing capacity.
- NiFi: observe processor throughput, back pressure, queue depth, and failure relationships. Adjust flow concurrency only after identifying whether ingestion, transformation, or publishing is the bottleneck.
- Pulsar: observe publish and consume rates, topic backlog, storage growth, and broker/storage health. Pulsar documents horizontal scaling, but that does not predict the capacity of this particular application.
- Flink: observe consumer progress, processing throughput, checkpoint completion and failures, task failures, and backpressure. Adjust parallelism only when measurements and topic partitioning support it.
Set operational alert thresholds from a measured baseline and service objectives; no numeric threshold can be inferred from the project architecture alone.
Quick Recap
Operational checklist before production
- Record exact NiFi, Pulsar, Flink, and connector versions, and keep each component’s configuration aligned with its corresponding release documentation.
- Verify the NiFi publishing mechanism and Flink connector in the actual secured environment, including credentials, TLS, serialization, and error behavior.
- Track Pulsar backlog and Flink consumer progress together so a growing queue can be distinguished from a slow or stopped processor.
- Monitor checkpoint health, job restarts, sink failures, NiFi queues, and Pulsar storage—not just whether each service is reachable.
- Exercise replay and recovery procedures, including the handling of duplicates and malformed events, before relying on the pipeline for critical workloads.
Product prices and availability are accurate as of the date/time indicated and are subject to change. Any price and availability information displayed on Amazon at the time of purchase will apply.

