Member since
01-15-2019
66
Posts
38
Kudos Received
2
Solutions
My Accepted Solutions
| Title | Views | Posted |
|---|---|---|
| 4321 | 07-20-2021 01:05 AM |
09-21-2026
06:38 AM
Most enterprises treat Salesforce as a system of record for the business, and a data platform (a Lakehouse, a warehouse, a set of streaming jobs) as the place where that business data gets joined, enriched, and analyzed. The gap between the two is deceptively simple to state — "get Salesforce changes into Kafka in real time" — and surprisingly easy to get wrong.
This post walks through the architecture of doing it properly: subscribing to Salesforce Change Data Capture (CDC), landing every change event in Kafka without loss, and feeding a medallion-style lakehouse downstream. It is design-focused; the code is illustrative, not a copy-paste tutorial.
The problem: there is no cable you can just plug-in
Salesforce publishes change events through its Pub/Sub API. This is not a REST endpoint you poll and not a Kafka broker you can point a consumer at. It is a gRPC service over HTTP/2, streaming Avro-encoded binary payloads, at api.pubsub.salesforce.com:7443. Its three RPCs matter to us:
Subscribe — a bidirectional stream. You send FetchRequest messages to request N events (flow control), and the server streams events back.
GetSchema — returns the Avro schema for a given schemaId.
GetTopic — metadata about a channel.
Kafka, on the other side, speaks its own binary protocol and knows nothing about gRPC or Salesforce Avro. Nothing in a standard data platform bridges these two natively. REST-based Salesforce connectors exist, but they poll objects — they are not a CDC event subscription and will miss the ordering, deletes, and low latency that CDC gives you. So the first architectural truth is:
Between the Salesforce Pub/Sub API and Kafka you must run a small bridge component that subscribes over gRPC, decodes Avro, and produces to Kafka.
The good news: Salesforce ships an official Pub/Sub API client library, so the bridge is closer to a wrapper than a from-scratch protocol implementation.
What a change event actually looks like
Before choosing where to run the bridge, it helps to see the payload. Every CDC event carries a stable envelope — the ChangeEventHeader — plus the record fields. A CREATE on an Account-like object decodes to roughly this:
{
"ChangeEventHeader": {
"entityName": "Account",
"recordIds": ["001XX000003ABCDEAG"],
"changeType": "CREATE",
"commitNumber": 1785576707554279428,
"commitTimestamp": 1785576707000,
"sequenceNumber": 1,
"transactionKey": "000047dd-d8ed-e5e5-1b48-498b18084963",
"changedFields": []
},
"Name": "Acme Corp",
"Industry": null,
"OwnerId": "005XX000001AbcdYAR"
}
Two things are worth internalizing:
The envelope is platform-wide and fixed. changeType (CREATE / UPDATE /
DELETE / UNDELETE), recordIds, commitNumber, commitTimestamp, sequenceNumber, transactionKey, and the changedFields bitmap are the same for every object in every org. Your downstream logic can be written once against this contract.
The payload fields are per-object. Which fields appear — and which custom
fields (*__c) show up — depends on which objects the org enabled CDC for. This is the one part you cannot know from documentation; it comes from the specific org's configuration.
That split (fixed envelope, variable payload) is what makes a generic, reusable pipeline possible.
The resume mechanism: replayId
Every event carries a replayId — an opaque byte bookmark marking the event's position in the stream. It is conceptually identical to a Kafka offset or a database CDC checkpoint (an SCN-style ordering clock, if you come from the Oracle world). When you subscribe you choose a ReplayPreset:
LATEST — only events from now on.
EARLIEST — from the start of the retention window.
CUSTOM — resume after a specific replayId. This is the one that prevents data loss.
The operational scenario is the whole ballgame. Your bridge stops — a deploy, a crash, a network blip — and Salesforce keeps producing events the whole time. If you restart with LATEST, everything produced during the outage is gone. If instead you persisted the replayId of the last event you successfully wrote to Kafka, you restart with CUSTOM from that bookmark and backfill the entire gap.
Two properties fall out of this, and you should state both plainly to stakeholders:
No loss. Resume-from-bookmark guarantees outage events are recovered.
At-least-once, not exactly-once. The Pub/Sub API delivers at least once. If the bridge writes to Kafka but dies before persisting the replayId, that event replays on restart. Duplicates at the boundary are normal and expected. You make the pipeline idempotent downstream by deduplicating on (recordId, commitNumber). The goal is "no loss," never "zero duplicates."
Where should the bridge run? Comparing the options
There are several defensible places to run the subscribe-decode-produce logic. Condensed, the realistic choices are:
OptionShapeVerdict
A
Bridge writes straight to the serving store (no Kafka)
Loses replay/buffer; couples ingestion to storage
B
Bridge → Kafka, then Flink downstream
Good, but "bridge" is undefined ops-wise
C
Bridge → Kafka → routing tool → store
Extra hop, little gain
D
Kafka Connect Source connector → Kafka → Flink
Recommended
E
Standalone microservice → Kafka → Flink
Works, but you own a service lifecycle
Option D — implement the bridge as a Kafka Connect Source connector — wins for one decisive reason: the framework already solves the hard part. In the Kafka Connect model, a SourceTask.poll() returns records each carrying a sourceOffset. Connect persists that offset automatically to an internal offsets topic and hands it back on restart. So the replayId ↔ offset bookkeeping — the thing you would otherwise write, test, and get subtly wrong — becomes map replayId onto the Connect source offset and let the framework do it.
You also inherit the rest of the Connect operational surface for free:
Horizontal scale via tasks. If you must subscribe to many orgs or many channels, Connect distributes tasks across workers. This matters enormously at fan-out (think hundreds of subscriptions).
Managed lifecycle, config, and monitoring through the Connect REST API and whatever streaming-management UI your platform ships.
Dead-letter handling for poison records without killing the task.
A standalone microservice (Option E) can do all of this too — but then you are building offset persistence, scaling, config management, and health checks. On a platform that already runs Kafka Connect (for example, the Streams Messaging stack in Cloudera Data Platform, which ships an Apache-based Kafka Connect managed by SMM), Option D is strictly less code you own.
Here is the architecture:
Enable the CDC feature in Salesforce
Click the [setup] icon in SFDC menu.
Search Change Data Capture / 変更データキャプチャ
Then configure what entity to sync.
(PoC) Use Python to subscribe the events and then save in Kafka
# pubsub_cdc_poc.py — Salesforce Pub/Sub API で CDC を購読する最小 PoC (CDP 非依存)
import io, queue
import avro.schema, avro.io
import grpc, requests
import pubsub_api_pb2 as pb2
import pubsub_api_pb2_grpc as pb2_grpc
# --- 1. OAuth ログイン (PoC=username-password / 生産は JWT bearer 推奨) ---
def sfdc_login(login_url, cid, csecret, user, pw):
r = requests.post(f"{login_url}/services/oauth2/token", data={
"grant_type": "password", "client_id": cid, "client_secret": csecret,
"username": user, "password": pw}) # pw = パスワード + security token
r.raise_for_status(); j = r.json()
org_id = j["id"].split("/")[-2] # id=".../<orgId>/<userId>" → orgId=tenantId
return j["access_token"], j["instance_url"], org_id
TOKEN, INSTANCE_URL, TENANT_ID = sfdc_login("https://login.salesforce.com", ...)
# --- 2. gRPC チャネル (443, TLS/HTTP2) + 認証メタデータ ---
channel = grpc.secure_channel("api.pubsub.salesforce.com:443", grpc.ssl_channel_credentials())
stub = pb2_grpc.PubSubStub(channel)
AUTH = (("accesstoken", TOKEN), ("instanceurl", INSTANCE_URL), ("tenantid", TENANT_ID))
TOPIC = "/data/活動対象__ChangeEvent" # ← 自定义对象の CDC 通道
# --- 3. Avro スキーマ取得 (schema_id ごとにキャッシュ) ---
_cache = {}
def avro_schema(schema_id):
if schema_id not in _cache:
info = stub.GetSchema(pb2.SchemaRequest(schema_id=schema_id), metadata=AUTH)
_cache[schema_id] = avro.schema.parse(info.schema_json)
return _cache[schema_id]
def decode(schema, payload):
return avro.io.DatumReader(schema).read(avro.io.BinaryDecoder(io.BytesIO(payload)))
# --- 4. Subscribe (双方向ストリーミング + フロー制御) ---
def run():
reqs = queue.Queue()
# 初回: LATEST から。再開時は replay_preset=CUSTOM + replay_id=<保存した replayId>
reqs.put(pb2.FetchRequest(topic_name=TOPIC,
replay_preset=pb2.ReplayPreset.LATEST, num_requested=10))
def req_iter():
while True: yield reqs.get()
for resp in stub.Subscribe(req_iter(), metadata=AUTH):
for ev in resp.events:
rec = decode(avro_schema(ev.event.schema_id), ev.event.payload)
h = rec["ChangeEventHeader"]
print(h["changeType"], h["entityName"], h["recordIds"], h["commitTimestamp"])
save_replay_id(ev.replay_id) # ★これを永続化 → 断点续接
if resp.pending_num_requested == 0: # 残枠が尽きたら補充 (フロー制御)
reqs.put(pb2.FetchRequest(topic_name=TOPIC, num_requested=10)) # 継続時は preset 不要
run()
I ran this on my MBP, it works.
(Production) Use Kafka Connect to subscribe the events and then save in Kafka
In a Product environment, we can create a Kafka Connect Source Connector
Source code example:
// SalesforcePubSubSourceConnector.java — 外殻。topic(entity)/org 単位で task 分割
public class SalesforcePubSubSourceConnector extends SourceConnector {
private Map<String,String> props;
@Override public void start(Map<String,String> p) { this.props = p; }
@Override public Class<? extends Task> taskClass() { return SalesforcePubSubSourceTask.class; }
@Override public List<Map<String,String>> taskConfigs(int maxTasks) {
List<String> topics = Arrays.asList(props.get("salesforce.topics").split(","));
// 購読トピック(=entity×org)を task へ配分。230 org はここで水平分割。
List<Map<String,String>> cfgs = new ArrayList<>();
for (List<String> g : ConnectorUtils.groupPartitions(topics, Math.min(maxTasks, topics.size()))) {
Map<String,String> c = new HashMap<>(props);
c.put("salesforce.topics", String.join(",", g));
cfgs.add(c);
}
return cfgs;
}
@Override public void stop() {}
@Override public ConfigDef config() { return CONFIG_DEF; }
@Override public String version() { return "0.1.0"; }
}
// SalesforcePubSubSourceTask.java — 内核は developerforce/pub-sub-api の Java stub をラップ
public class SalesforcePubSubSourceTask extends SourceTask {
private PubSubClient client; // gRPC channel(443)+OAuth(JWT) のラッパー
private String kafkaTopic;
private final BlockingQueue<SourceRecord> buffer = new LinkedBlockingQueue<>();
@Override public void start(Map<String,String> props) {
kafkaTopic = props.get("kafka.topic"); // 例: sfdc.toyotec.katsudo
client = new PubSubClient(props);
for (String t : props.get("salesforce.topics").split(",")) {
// partition に org を入れておくと 230 org でも offset が衝突しない
Map<String,Object> partition = Map.of("sfdcTopic", t, "org", props.get("org.id"));
// ★ 再開: 前回 replayId を offset storage から復元 → CUSTOM で続き
Map<String,Object> off = context.offsetStorageReader().offset(partition);
byte[] replayId = off == null ? null : ((ByteBuffer) off.get("replayId")).array();
ReplayPreset preset = replayId == null ? ReplayPreset.LATEST : ReplayPreset.CUSTOM;
client.subscribe(t, preset, replayId, this::onEvent); // 非同期ストリーム
}
}
// gRPC で届いた変更イベント → SourceRecord にして buffer へ
private void onEvent(String sfdcTopic, String org, byte[] replayId, GenericRecord ev) {
Map<String,Object> partition = Map.of("sfdcTopic", sfdcTopic, "org", org);
Map<String,Object> offset = Map.of("replayId", ByteBuffer.wrap(replayId)); // ★肝
buffer.offer(new SourceRecord(partition, offset, kafkaTopic,
Schema.STRING_SCHEMA, recordId(ev), // key = org+PK 推奨(保序)
/* value schema */ null, ev.toString())); // 実際は Avro/Struct 変換
}
@Override public List<SourceRecord> poll() throws InterruptedException {
List<SourceRecord> batch = new ArrayList<>();
SourceRecord first = buffer.poll(1, TimeUnit.SECONDS);
if (first != null) { batch.add(first); buffer.drainTo(batch, 500); }
return batch; // ← framework が offset(replayId) を自動 commit
}
@Override public void stop() { if (client != null) client.close(); }
@Override public String version() { return "0.1.0"; }
}
Downstream: one CDC paradigm, reused
Once change events are in Kafka, the downstream is not Salesforce-specific anymore — it is just CDC, the same shape you already handle for database sources. A Flink job reads the topic, inspects ChangeEventHeader.changeType, and routes:
Silver (current state) — upsert/delete by recordId into a mutable store such as Kudu. This is the "what does the record look like now" table.
Bronze (history) — append every event to Iceberg for full audit, replay, and time travel.
Gold — modeled, business-ready tables built from Silver/Bronze.
The strategic payoff: if you already run an Oracle/DB CDC pipeline into the same medallion layout, Salesforce CDC folds into the identical Flink logic rather than becoming a second, parallel implementation. You unify on one CDC engine and avoid re-implementing history, logical deletes, and metadata handling twice.
The failure modes worth designing for up front
An honest architecture names what can go wrong. Three things dominate:
Schema evolution. The envelope is stable, but payloads gain and lose fields as admins change objects. Register schemas and enforce backward- compatible evolution so added fields are absorbed harmlessly. This is the most common cause of "the job suddenly can't parse an event."
Poison records. Route un-parseable events to a dead-letter topic instead of crashing the task; fix and replay later. Because the raw event is already durably in Kafka (and optionally landed raw in Bronze), a transform you can't handle today is never data loss — you reprocess from the offset once the logic is fixed. Kafka + Bronze is your replayable source of truth.
Authentication for unattended runs. Interactive OAuth is fine for a laptop proof-of-concept, but a long-running connector needs the JWT Bearer flow (a connected app with a certificate) so it can mint and refresh tokens headlessly. Decide this before implementation — it has a lead time (certificates, admin setup).
And one environmental prerequisite that stalls more projects than any code bug: network egress. Your platform typically lives in a private network; it needs an allowed outbound TLS path to api.pubsub.salesforce.com:7443. Confirm it early.
Proving it before you build it
You do not need the full connector to de-risk the design. A minimal bridge — a ~200-line script that authenticates, calls Subscribe, decodes Avro with the schema fetched via GetSchema, produces JSON to a local single-node Kafka, and writes the last replayId to a state file — is enough to validate the two things that actually carry risk:
No loss under steady state: create a known batch of records, confirm every recordId lands and the changeType counts match.
Resume correctness: start the bridge, kill it, make several changes while it is down, restart, and confirm the downtime events all arrive via CUSTOM replay — with boundary duplicates converging idempotently on (recordId, commitNumber).
That local, self-contained test reproduces the production design's core behavior (what Kafka Connect will later automate) and turns "we think this works" into "we watched it work." It is also the cheapest possible artifact to show a skeptical stakeholder.
Takeaways
The Salesforce Pub/Sub API is gRPC + Avro; you need a bridge to reach Kafka. There is no native connector, and REST polling is not CDC.
Implement the bridge as a Kafka Connect Source connector so the framework handles offset persistence, scaling, and lifecycle — map replayId onto the Connect source offset.
Design for at-least-once: guarantee no loss via replay, and make the downstream idempotent on (recordId, commitNumber).
Keep the raw event in Kafka (and Bronze) so any transform failure is reprocessable, never lost.
Settle schema evolution, dead-lettering, JWT auth, and network egress before you write the connector — those, not the gRPC plumbing, are where projects actually stall.
The plumbing is interesting, but the durable lesson is the same one that applies to every CDC integration: make ingestion lossless and replayable, push idempotency downstream, and reuse one change-data paradigm across all your sources. Salesforce just happens to speak gRPC on the way in.
... View more
06-10-2026
01:06 AM
1 Kudo
In the GenAI era, I’ve really enjoyed how tools like LLMs can boost productivity in day-to-day operations. However, I found that NiFi’s Web UI doesn’t integrate very naturally with GenAI workflows. I also explored options like MCP Server to bridge that gap, but it requires additional setup and a specific skillset.
Recently, I discovered a much simpler approach—leveraging NiFi’s REST APIs directly. Without any MCP setup, we can use LLMs as an operations co-pilot to extract, audit, and document NiFi flows and Ranger policies.
What Is Vibe Coding for NiFi Development?
Vibe coding means you describe your intent in natural language and let AI write the code. Applied to CDP NiFi flow development, it looks like this:
You say: "Show me all Process Groups and their processor counts"
AI does: Generates the correct curl commands, calls the NiFi REST API, parses the JSON, and returns a clean summary table
As a NiFi flow developer, you spend most of your time designing data pipelines -- not memorizing API paths. By pairing an AI assistant (like Claude) with the NiFi & Ranger REST APIs, you can:
Develop your flow with description language
Explore your flow topology and processor configurations instantly
Inspect Controller Services, connections, and versioning without clicking through the UI
Review Ranger security policies that govern your flows
Document your flow architecture for team handoff or migration
Iterate faster -- ask follow-up questions, drill into specific Process Groups, compare environments
This article provides ready-to-use prompt templates to get started.
Prerequisites
A running CDP environment with NiFi and Ranger services
Workload credentials (username/password) with API access
Network access to the CDP management endpoints (direct or via proxy)
An AI assistant (Claude, ChatGPT, etc.) that can execute or generate curl commands
Step 1: Set Up Environment Variables
Store your endpoints and credentials as environment variables. Never hardcode credentials in scripts or documents.
# --- Credentials (store in a separate file, e.g. .cdp-creds, and gitignore it) ---
WORKLOAD_USERNAME=<your-workload-username>
WORKLOAD_PASSWORD=<your-workload-password>
# --- NiFi Endpoints ---
# Find these in CDP > Data Hub > your NiFi cluster > Endpoints tab
NIFI_BASE_HOST=<nifi-management-host>.cloudera.site
NIFI_BASE_URL=https://$NIFI_BASE_HOST/<cluster-name>/cdp-proxy-api/nifi-app/nifi-api/
# --- Ranger Endpoints ---
# Find these in CDP > Data Lake > Endpoints tab
DATA_LAKE_GATEWAY=<datalake-gateway-host>.cloudera.site
RANGER_BASE_URL=https://$DATA_LAKE_GATEWAY/<datalake-name>/cdp-proxy-api/ranger/
Tip: Create a .cdp-creds file with chmod 600 , add it to .gitignore , and source it before running commands.
Step 2: Verify Connectivity
CDP in Public Subnet (Direct Access)
# Test NiFi connectivity
curl -s "${NIFI_BASE_URL}system-diagnostics" \
-u "$WORKLOAD_USERNAME:$WORKLOAD_PASSWORD" \
-o /dev/null -w "NiFi: HTTP %{http_code}\n"
# Test Ranger connectivity
curl -s "${RANGER_BASE_URL}service/public/v2/api/service" \
-u "$WORKLOAD_USERNAME:$WORKLOAD_PASSWORD" \
-o /dev/null -w "Ranger: HTTP %{http_code}\n"
# Expected output:
# NiFi: HTTP 200
# Ranger: HTTP 200
CDP in Private Subnet (via SOCKS5 Proxy)
If your CDP cluster is in a private subnet, use SSH dynamic forwarding through a jump host:
# Architecture: Local Machine --> Bastion/Jump Host --> CDP Cluster Internal Network
# Step 1: Set up SOCKS5 proxy via SSH tunnel to your jump host
ssh -D 1084 -N -f user@jump-host
# Step 2: Use socks5h:// (the 'h' means remote DNS resolution)
curl -s --proxy socks5h://127.0.0.1:1084 \
"${NIFI_BASE_URL}system-diagnostics" \
-u "$WORKLOAD_USERNAME:$WORKLOAD_PASSWORD" \
-o /dev/null -w "NiFi: HTTP %{http_code}\n"
# socks5h:// resolves hostnames on the proxy side (jump host)
# No need for --resolve hacks or /etc/hosts modifications
Step 3: Start Vibe Coding--Prompt Templates
Below are ready-to-use prompt templates. Copy one to your AI assistant, replace {{ }} placeholders, and let it do the heavy lifting.
Prompt 1: Explore Flow Topology
"I just joined the project -- give me a map of what's running."
You are a NiFi flow developer assistant. Help me understand the current flow
topology via the NiFi REST API.
## Connection Info
- NiFi API Base URL: {{ NIFI_BASE_URL }}
- Authentication: Basic Auth
- Username/Password: {{ WORKLOAD_USERNAME }} / {{ WORKLOAD_PASSWORD }}
## Tasks
Use curl to complete the following and output results as structured Markdown:
1. Get all Process Group list (name, ID, status)
GET /process-groups/root/process-groups
2. Get the complete flow configuration of a specific Process Group
GET /process-groups/{id}/flow
3. List all Controller Services (type, name, status, properties)
GET /flow/controller/controller-services
4. Export a specific Process Group as template JSON
POST /process-groups/{id}/templates/export
## Output Format
- Summarize Processor list in a Markdown table (name, type, status, scheduling strategy)
- Save complete JSON configuration in code blocks
- Flag sensitive properties that require manual review (passwords, keys, etc.)
Prompt 2: Inspect Flow Versioning & Registry
"Which flows are version-controlled? What changed in the last release?"
You are a NiFi flow developer assistant. Help me inspect Registry and version
control configuration via the NiFi REST API.
## Connection Info
- NiFi API Base URL: {{ NIFI_BASE_URL }}
- Authentication: Basic Auth ({{ WORKLOAD_USERNAME }})
## Tasks
1. List all Registry Client configurations
GET /controller/registry-clients
2. View versioned Process Group list
GET /flow/registries/{registryId}/buckets/{bucketId}/flows
3. Get a specific version's Flow snapshot
GET /buckets/{bucketId}/flows/{flowId}/versions/{versionNumber}
## Output
- Registry inventory (name, URL, type)
- Version history table for each Flow
Prompt 3: Review Ranger Security Policies
"Who has access to what? Are there any overly permissive policies?"
You are a NiFi flow developer assistant. Help me review the Ranger security
policies that govern my NiFi flows.
## Connection Info
- Ranger API Base URL: {{ RANGER_BASE_URL }}
- Authentication: Basic Auth ({{ WORKLOAD_USERNAME }})
## Tasks
1. Get all Service (plugin) list
GET /service/public/v2/api/service
2. Get all Policies for a specific Service
GET /service/public/v2/api/policy?serviceName={{ SERVICE_NAME }}
3. Get all users and user groups
GET /service/xusers/users
GET /service/xusers/groups
4. Get role list
GET /service/roles/roles
## Output Format
- Policy summary table: Policy Name | Resource Path | Allowed Users/Groups | Permissions | Enabled
- User-Role mapping table
- Flag high-privilege policies (containing * wildcard or admin permissions)
Prompt 4: Audit NiFi-Specific Ranger Policies
"My processor is getting 'access denied' -- what Ranger policies apply to my flow?"
You are a NiFi flow developer assistant. Help me find and audit the Ranger
policies specific to NiFi.
## Connection Info
- Ranger API Base URL: {{ RANGER_BASE_URL }}
- Authentication: Basic Auth ({{ WORKLOAD_USERNAME }})
## Tasks
1. Find NiFi Service definition
GET /service/public/v2/api/service?serviceType=nifi
2. Get all policies under that Service
GET /service/public/v2/api/policy?serviceName={{ NIFI_SERVICE_NAME }}
3. Export policies as JSON file (for backup or migration)
## Output
- NiFi resource path permission inventory
- Access control policies for each Process Group
- Flag entries that conflict with or duplicate NiFi built-in policies
Prompt 5: Full Flow Documentation for Migration
"We're migrating to a new environment -- generate the complete inventory."
You are a NiFi flow developer assistant. Generate a complete configuration
inventory of NiFi flows and Ranger policies for documentation and migration.
## Connection Info
- NiFi API Base URL: {{ NIFI_BASE_URL }}
- Ranger API Base URL: {{ RANGER_BASE_URL }}
- Authentication: Basic Auth ({{ WORKLOAD_USERNAME }})
## Execution Steps (in order)
### Step 1: NiFi Flow Inventory
- Get all Process Groups (including hierarchy)
- Count Processor quantities and type distribution in each Group
### Step 2: NiFi Controller Services
- List all Controller Services (type, name, status)
- Flag critical services such as DBCPConnectionPool, SSLContextService, etc.
### Step 3: Ranger Policy Inventory
- Get all NiFi-related Services and Policies
- Output Policy permission matrix
### Step 4: Summary Report
Generate a Markdown report containing:
- NiFi flow architecture overview (Mermaid diagram)
- Controller Service configuration table
- Ranger permission matrix
- Discovered issues or recommendations
## Output Files
- nifi-flow-inventory.md
- ranger-policy-matrix.md
- full-config-report.md
API Quick Reference
NiFi REST API
Operation Method Path
Get root Process Group
GET
/process-groups/root
List child Process Groups
GET
/process-groups/{id}/process-groups
Get Flow details
GET
/process-groups/{id}/flow
List Processors
GET
/process-groups/{id}/processors
List Connections
GET
/process-groups/{id}/connections
List Controller Services
GET
/flow/controller/controller-services
Get system diagnostics
GET
/system-diagnostics
Get cluster status
GET
/controller/cluster
Token authentication
POST
/access/token
Ranger REST API
Operation Method Path
List all Services
GET
/service/public/v2/api/service
List all Policies
GET
/service/public/v2/api/policy
Query Policies by Service
GET
/service/public/v2/api/policy?serviceName={name}
Create Policy
POST
/service/public/v2/api/policy
List users
GET
/service/xusers/users
List user groups
GET
/service/xusers/groups
List roles
GET
/service/roles/roles
Service type definitions
GET
/service/public/v2/api/servicedef
The Vibe Coding Workflow
Set up environment variables and credentials ( .cdp-creds file)
Verify connectivity with a quick curl test against /system-diagnostics
Describe what you want to know or do -- in natural language
AI generates the correct curl commands targeting NiFi/Ranger APIs
AI parses the JSON responses and produces structured reports
Iterate -- ask follow-up questions, drill into specific flows, compare environments
Real-World Example: Building a CSV-to-S3 Ingestion Flow with AI
Here's what vibe coding NiFi development looks like in practice. Instead of dragging processors in the UI and manually configuring each one, you have a conversation with AI:
The Conversation
You:
I need to build a flow that picks up CSV files from a local directory, validates the schema, converts them to Parquet, and writes to S3. Set it up in my "Data-Ingestion" Process Group.
AI responds with the exact REST API calls:
# Step 1: Find the target Process Group
curl -s "${NIFI_BASE_URL}process-groups/root/process-groups" \
-u "$WORKLOAD_USERNAME:$WORKLOAD_PASSWORD" | jq '.processGroups[] | {name: .component.name, id: .id}'
# Step 2: Create a GetFile processor
curl -s -X POST "${NIFI_BASE_URL}process-groups/<pg-id>/processors" \
-u "$WORKLOAD_USERNAME:$WORKLOAD_PASSWORD" \
-H "Content-Type: application/json" \
-d '{
"revision": {"version": 0},
"component": {
"type": "org.apache.nifi.processors.standard.GetFile",
"name": "Pick up CSV files",
"config": {
"properties": {
"Input Directory": "/data/incoming/csv",
"File Filter": "[^\\.].*\\.csv",
"Keep Source File": "false"
}
}
}
}'
# Step 3: Create a ValidateRecord processor for schema validation
curl -s -X POST "${NIFI_BASE_URL}process-groups/<pg-id>/processors" \
-u "$WORKLOAD_USERNAME:$WORKLOAD_PASSWORD" \
-H "Content-Type: application/json" \
-d '{
"revision": {"version": 0},
"component": {
"type": "org.apache.nifi.processors.standard.ValidateRecord",
"name": "Validate CSV Schema",
"config": {
"properties": {
"Record Reader": "<csv-reader-controller-service-id>",
"Record Writer": "<csv-writer-controller-service-id>"
}
}
}
}'
# Step 4: Create a ConvertRecord processor (CSV -> Parquet)
curl -s -X POST "${NIFI_BASE_URL}process-groups/<pg-id>/processors" \
-u "$WORKLOAD_USERNAME:$WORKLOAD_PASSWORD" \
-H "Content-Type: application/json" \
-d '{
"revision": {"version": 0},
"component": {
"type": "org.apache.nifi.processors.standard.ConvertRecord",
"name": "Convert CSV to Parquet",
"config": {
"properties": {
"Record Reader": "<csv-reader-controller-service-id>",
"Record Writer": "<parquet-writer-controller-service-id>"
}
}
}
}'
# Step 5: Create a PutS3Object processor
curl -s -X POST "${NIFI_BASE_URL}process-groups/<pg-id>/processors" \
-u "$WORKLOAD_USERNAME:$WORKLOAD_PASSWORD" \
-H "Content-Type: application/json" \
-d '{
"revision": {"version": 0},
"component": {
"type": "org.apache.nifi.processors.aws.s3.PutS3Object",
"name": "Write Parquet to S3",
"config": {
"properties": {
"Bucket": "my-data-lake-bucket",
"Object Key": "ingested/${now():format('\''yyyy/MM/dd'\'')}/${filename}.parquet"
}
}
}
}'
# Step 6: Connect the processors in sequence
# GetFile -> ValidateRecord -> ConvertRecord -> PutS3Object
curl -s -X POST "${NIFI_BASE_URL}process-groups/<pg-id>/connections" \
-u "$WORKLOAD_USERNAME:$WORKLOAD_PASSWORD" \
-H "Content-Type: application/json" \
-d '{
"revision": {"version": 0},
"component": {
"source": {"id": "<getfile-id>", "type": "PROCESSOR"},
"destination": {"id": "<validate-id>", "type": "PROCESSOR"},
"selectedRelationships": ["success"]
}
}'
# ... (AI generates the remaining connections)
Then You Iterate
You: "Add error handling -- route validation failures to a separate PutFile that writes to /data/rejected/ "
AI generates the additional processor and connection API calls.
You: "What Controller Services do I need for the CSV reader and Parquet writer?"
AI creates the CSVReader and ParquetRecordSetWriter Controller Services:
# Create CSVReader Controller Service
curl -s -X POST "${NIFI_BASE_URL}process-groups/<pg-id>/controller-services" \
-u "$WORKLOAD_USERNAME:$WORKLOAD_PASSWORD" \
-H "Content-Type: application/json" \
-d '{
"revision": {"version": 0},
"component": {
"type": "org.apache.nifi.csv.CSVReader",
"name": "CSV Reader",
"properties": {
"Schema Access Strategy": "Use String Fields From Header",
"Date Format": "yyyy-MM-dd",
"Timestamp Format": "yyyy-MM-dd HH:mm:ss"
}
}
}'
# Create ParquetRecordSetWriter Controller Service
curl -s -X POST "${NIFI_BASE_URL}process-groups/<pg-id>/controller-services" \
-u "$WORKLOAD_USERNAME:$WORKLOAD_PASSWORD" \
-H "Content-Type: application/json" \
-d '{
"revision": {"version": 0},
"component": {
"type": "org.apache.nifi.parquet.ParquetRecordSetWriter",
"name": "Parquet Writer",
"properties": {
"Schema Access Strategy": "Inherit Record Schema",
"Compression Type": "SNAPPY"
}
}
}'
# Enable the Controller Services
curl -s -X PUT "${NIFI_BASE_URL}controller-services/<service-id>/run-status" \
-u "$WORKLOAD_USERNAME:$WORKLOAD_PASSWORD" \
-H "Content-Type: application/json" \
-d '{"revision": {"version": 1}, "state": "ENABLED"}'
You: "Now start the entire Process Group."
# Start all processors in the Process Group
curl -s -X PUT "${NIFI_BASE_URL}flow/process-groups/<pg-id>" \
-u "$WORKLOAD_USERNAME:$WORKLOAD_PASSWORD" \
-H "Content-Type: application/json" \
-d '{"id": "<pg-id>", "state": "RUNNING"}'
What Just Happened?
In a few minutes of conversation, you built a complete ingestion pipeline: No UI clicking. No API doc searching. No API doc searching. Just describing, reviewing, and iterating.
Security Best Practices
Never commit credentials to Git. Use environment variables loaded from a gitignored file.
Use chmod 600 on credential files to restrict read access.
Rotate workload passwords regularly through CDP Management Console.
Review AI output before acting on any write operations -- always treat AI suggestions as recommendations, not commands.
Use read-only operations (GET) for exploration. Only use write operations (POST/PUT/DELETE) after careful review.
Conclusion
Vibe coding flips the NiFi development workflow: instead of navigating the UI or looking up REST API docs, you describe your intent and let AI handle the plumbing. Whether you're onboarding to an existing project, debugging access control issues, or preparing a migration inventory, these prompt templates let you stay in the flow -- pun intended.
Applicable to: Cloudera Data Platform (CDP) Public Cloud -- Data Hub NiFi clusters and Data Lake Ranger services.
... View more
Labels:
12-17-2025
04:14 AM
@zhouweibo Hi Weibo, I build an environment by myself, but I can't reproduce your error. I created a table in this way: -- Create database if not exists
CREATE DATABASE IF NOT EXISTS upidb;
-- Create table with correct syntax
CREATE TABLE IF NOT EXISTS upidb.gscs_tbl_fultrans_2_db (
trace_num STRING COMMENT 'Trace number',
acq_ins_cde STRING COMMENT 'Acquiring institution code',
fwd_ins_cde STRING COMMENT 'Forwarding institution code',
acq_trans_cde STRING COMMENT 'Acquiring transaction code',
iss_trans_cde STRING COMMENT 'Issuing transaction code',
pri_acct_num STRING COMMENT 'Primary account number',
resv_fld2 STRING COMMENT 'Reserved field 2 - contains encoded data',
sett_dt STRING COMMENT 'Settlement date in YYYYMMDD format'
)
COMMENT 'Full transaction table for settlement date partition'
STORED AS PARQUET;
-- Insert sample data with GBK encoded characters to reproduce the encoding issue
INSERT INTO upidb.gscs_tbl_fultrans_2_db VALUES
(
'123456789012',
'ACQ001',
'FWD001',
'TRANS001',
'TRANS002',
'1234567890123456',
'ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789测试中文字符ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZ',
'20231201'
); Then ran the SELECT sql, but Can't reproduce error: SELECT trace_num
,acq_ins_cde
,fwd_ins_cde
,acq_trans_cde
,iss_trans_cde
,pri_acct_num
,trim(SUBSTR(resv_fld2,111,2))
,trim(decode(SUBSTR(encode(resv_fld2,'gbk'),113,15),'gbk'))
,trim(decode(SUBSTR(encode(resv_fld2,'gbk'),128,40),'gbk'))
,trim(decode(SUBSTR(encode(resv_fld2,'gbk'),168,40),'gbk'))
from upidb.gscs_tbl_fultrans_2_db
where sett_dt = '20231201' Can you please share your DDL?
... View more
10-03-2025
01:47 AM
Since the previous link (hortonworks.com) has expired, please refer to the updated links below: Change Data Capture (CDC) with Apache NiFi – Part 1 of 3 https://community.cloudera.com/t5/Community-Articles/Change-Data-Capture-CDC-with-Apache-NiFi-Part-1-of-3/ta-p/246623 Change Data Capture (CDC) with Apache NiFi – Part 2 of 3 https://community.cloudera.com/t5/Community-Articles/Change-Data-Capture-CDC-with-Apache-NiFi-Part-2-of-3/ta-p/246519 Change Data Capture (CDC) with Apache NiFi – Part 3 of 3 https://community.cloudera.com/t5/Community-Articles/Change-Data-Capture-CDC-with-Apache-NiFi-Part-3-of-3/ta-p/246482
... View more
10-03-2025
01:46 AM
Since the previous link (hortonworks.com) has expired, please refer to the updated links below: Change Data Capture (CDC) with Apache NiFi – Part 1 of 3 https://community.cloudera.com/t5/Community-Articles/Change-Data-Capture-CDC-with-Apache-NiFi-Part-1-of-3/ta-p/246623 Change Data Capture (CDC) with Apache NiFi – Part 2 of 3 https://community.cloudera.com/t5/Community-Articles/Change-Data-Capture-CDC-with-Apache-NiFi-Part-2-of-3/ta-p/246519 Change Data Capture (CDC) with Apache NiFi – Part 3 of 3 https://community.cloudera.com/t5/Community-Articles/Change-Data-Capture-CDC-with-Apache-NiFi-Part-3-of-3/ta-p/246482
... View more
10-03-2025
01:46 AM
Since the previous link (hortonworks.com) has expired, please refer to the updated links below: Change Data Capture (CDC) with Apache NiFi – Part 1 of 3 https://community.cloudera.com/t5/Community-Articles/Change-Data-Capture-CDC-with-Apache-NiFi-Part-1-of-3/ta-p/246623 Change Data Capture (CDC) with Apache NiFi – Part 2 of 3 https://community.cloudera.com/t5/Community-Articles/Change-Data-Capture-CDC-with-Apache-NiFi-Part-2-of-3/ta-p/246519 Change Data Capture (CDC) with Apache NiFi – Part 3 of 3 https://community.cloudera.com/t5/Community-Articles/Change-Data-Capture-CDC-with-Apache-NiFi-Part-3-of-3/ta-p/246482
... View more
08-19-2025
01:05 AM
Several keys needed to be added: This is an example of the properties we used in KConnect in DH ---------------------------- 1- producer.override.sasl.jaas.config org.apache.kafka.common.security.plain.PlainLoginModule required username="<your-workload-name>" password="<password>"; 2- producer.override.security.protocol SASL_SSL 3- producer.override.sasl.mechanism PLAIN ----------------------------
... View more
08-18-2025
08:52 AM
H i@shubham_rai Do you have a chance to try the Custom Service on a CDP Base (on-premises) version? If you run it on CDP On-Premises, do you get the same error message?
... View more
08-17-2025
10:28 PM
@cnelson2 This is really helpful! Thanks!
... View more
05-23-2025
06:19 AM
1 Kudo
This guide provides a step-by-step approach to extracting data from SAP S/4HANA via OData APIs, processing it using Apache NiFi in Cloudera Data Platform (CDP), and storing it in an Iceberg-based Lakehouse for analytics and AI workloads. 1. Introduction 1.1 Why Move SAP S/4HANA Data to a Lakehouse? SAP S/4HANA is a powerful ERP system designed for transactional processing, but it faces limitations when used for analytics, AI, and large-scale reporting: Performance Impact: Running complex analytical queries directly on SAP can degrade system performance. Limited Scalability: SAP systems are not optimized for big data workloads (e.g., petabyte-scale analytics). High Licensing Costs: Extracting and replicating SAP data for analytics can be expensive if done inefficiently. Lack of Flexibility: SAP’s data model is rigid, making it difficult to integrate with modern AI/ML tools. A Lakehouse architecture (built on Apache Iceberg in CDP) solves these challenges by: Decoupling analytics from SAP – Reduce operational load on SAP while enabling scalable analytics. Supporting structured & unstructured data – Unlike SAP’s tabular model, a Lakehouse can store JSON, text, and IoT data. Enabling ACID compliance – Iceberg ensures transactional integrity (critical for financial and inventory data). Reducing costs – Store historical SAP data in cheaper object storage (S3, ADLS) rather than expensive SAP HANA storage. 1.2 Why Use OData API for SAP Data Extraction? SAP provides several data extraction methods, but OData (Open Data Protocol) is one of the most efficient for real-time replication: Method Pros Cons Best For OData API Real-time, RESTful, easy to use Requires pagination handling Incremental, near-real-time syncs SAP BW/Extractors SAP-native, optimized for BW Complex setup, not real-time Legacy SAP BW integrations Database Logging (CDC) Low latency, captures all changes High SAP system overhead Mission-critical CDC use cases SAP SLT (Trigger-based) Real-time, no coding needed Expensive, SAP-specific Large-scale SAP replication Why OData wins for Lakehouse ingestion? REST-based – Works seamlessly with NiFi’s InvokeHTTP processor. Supports filtering ($filter) – Enables incremental extraction (e.g., modified_date gt ‘2024-01-01’). JSON/XML output – Easy to parse and transform in NiFi before loading into Iceberg. 1.3 Why Apache NiFi in Cloudera Data Platform (CDP)? NiFi is the ideal tool for orchestrating SAP-to-Lakehouse pipelines because: Low-Code UI: Drag-and-drop processors simplify pipeline development (vs. writing custom Spark/PySpark code). Built-in SAP Connectors: Use InvokeHTTP for SAP S/4 HANA OData for deeper integrations. Scalability & Fault Tolerance: Backpressure handling – Prevents SAP API overload. Automatic retries – If SAP API fails, NiFi retries without data loss. 2. Prerequisites Before building the SAP S/4HANA → NiFi → Iceberg pipeline, ensure the following components and access rights are in place. Cloudera Data Platform (CDP) with: Apache NiFi (for data ingestion) Apache Iceberg (as the Lakehouse table format) Storage: HDFS or S3 (via Cloudera SDX) SAP S/4HANA access with OData API permissions T-Code SEGW: Confirm OData services are exposed (e.g., API_MATERIAL_SRV). Permissions: SAP User Role: Must include S_ODATA and S_RFC authorizations. Whitelist NiFi IP if SAP has network restrictions. Test OData Endpoints curl -u "USER:PASS" "https://sap-odata.example.com:443/sap/opu/odata/sap/API_SALES_ORDER_SRV/A_SalesOrder?$top=2" Validate: Pagination ($skip, $top). Filtering ($filter=LastModified gt '2025-05-01'). Basic knowledge of NiFi flows, SQL, and Iceberg 3. Architecture Overview Data movement: SAP S/4HANA (OData API) → Apache NiFi (CDP) → Iceberg Tables (Lakehouse) → Analytics (Spark, Impala, Hive) Archtecture Overview : 4. Step-by-Step Implementation Step 1: Identify SAP OData Endpoints SAP provides OData services for tables like: MaterialMaster (MM) SalesOrders (SD) FinancialDocuments (FI) Example endpoint: https://<SAP_HOST>:<PORT>/sap/opu/odata/sap/API_SALES_ORDER_SRV/A_SalesOrder?$top=2 Step 2: Configure NiFi to Extract SAP Data Use InvokeHTTP processor to call SAP OData API. Configure authentication (Basic Auth). Handle pagination ($skip & $top parameters). To get the JSON response, I added Accept=application/json Property. Parse JSON responses using EvaluateJsonPath or JoltTransformJSON. Step 3: Transform Data in NiFi Filter & clean data using: ReplaceText (for SAP-specific formatting) QueryRecord (to convert JSON to Parquet/AVRO) Enrich data (e.g., join with reference tables). Check the Data using Provinance : Step 4: Load into Iceberg Lakehouse Use PutIceberg processor (NiFi 1.23+) to write directly to Iceberg. Alternative Option: Write to HDFS/S3 as Parquet, then use Spark SQL to load into Iceberg CREATE TABLE iceberg_db.sap_materials (
material_id STRING,
material_name STRING,
created_date TIMESTAMP
)
STORED AS ICEBERG; 5. Conclusion By leveraging Cloudera’s CDP, NiFi, and Iceberg, organizations can efficiently move SAP data into a modern Lakehouse, enabling real-time analytics, ML, and reporting without impacting SAP performance. Next Steps Explore Cloudera Machine Learning (CML) for SAP data analytics.
... View more
Labels: