- Support
- Integrations
- Streaming
Streaming
Oracle Cloud integration · 30 node(s).
00Overview
Publish and consume events on Oracle Cloud's Kafka-compatible Streaming service straight from a flow — put messages onto a stream, read them back with partition or consumer-group cursors, and commit or reset where a group resumes. Manage the surrounding resources too: create and list streams, stream pools and connect harnesses, move them between compartments, and update their tags. Long-running create, update and delete operations return a work request you can poll from a later step to confirm the change landed.
Every field below is exactly what you see in the Flomation editor. Fields marked ● live picker let you choose from a list pulled live from your account — no IDs to look up.
01Connecting Streaming
- In every Streaming node, set Authentication to either Connect Oracle Cloud — pick an already-connected Oracle Cloud account in the Oracle Cloud connection field and skip the rest of this list — or API signing key (advanced) to enter signing-key credentials by hand.
- For the API signing key method, sign in to the Oracle Cloud Console and open Profile → My profile → API keys, choose Add API key, and download the generated private key.
- Copy the values from the configuration-file preview Oracle shows you into the node: Tenancy OCID, User OCID, Region (e.g.
uk-london-1) and Key Fingerprint. Set Compartment OCID to the compartment that holds — or will hold — your streams (Identity → Compartments). - In Flomation, add the downloaded private key (the full PEM block) as an environment secret (e.g.
streaming_secret) and pick it in the node's Private Key (PEM) field; if you set a passphrase when generating the key, store that as a secret too and select it in Private Key Passphrase.
| Field | Type | Details | |
|---|---|---|---|
| Authentication | string | Connect Oracle Cloud, API signing key (advanced) | |
| Oracle Cloud connection | credential | Pick a connected Oracle Cloud account | |
| Region | string | e.g. uk-london-1 | |
| Private Key (PEM) | secret | The API signing private key — full PEM, incl. BEGIN/END lines | |
| Private Key Passphrase | secret | Only if the key is encrypted (optional) | |
| Tenancy OCID | string | ocid1.tenancy.oc1..aaaa… | |
| User OCID | string | ocid1.user.oc1..aaaa… | |
| Key Fingerprint | string | aa:bb:cc:… fingerprint of the uploaded API key |
Pick an Environment on your flow (Flow Settings → Environment) so the secret resolves. Secret fields never show the value — they reference ${secrets.your_secret}.
02Connect
OCI Streaming: Move Connect Harness to Compartment
oracle/streaming/connect_harness_change_compartment · Action
Move a connect harness into a different compartment. Supply the connect harness OCID and the destination compartment OCID. The move is asynchronous and returns a work request OCID to poll.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Connect Harness OCID | string | Required | ocid1.connectharness.oc1..aaaa… |
| Destination Compartment OCID | string | Required | ocid1.compartment.oc1..aaaa… (where the connect harness should move to) |
Returns: tool_result, id, destination_compartment_id, work_request_id, success, error
OCI Streaming: Create Connect Harness
oracle/streaming/connect_harness_create · Action
Create a connect harness for Kafka Connect support. Give it a name and a compartment. Returns the harness in a CREATING state — poll Get Connect Harness until ACTIVE.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | Required | ocid1.compartment.oc1..aaaa… (holds the connect harness) |
| Connect Harness Name | string | Required | e.g. JDBCConnector |
| Freeform Tags (JSON) | string | {"env":"prod"} (optional) | |
| Defined Tags (JSON) | string | {"Ops":{"env":"prod"}} (optional) |
Returns: tool_result, connect_harness, id, lifecycle_state, work_request_id, success, error
OCI Streaming: Delete Connect Harness
oracle/streaming/connect_harness_delete · Action
Delete a connect harness by OCID. Asynchronous — returns a work request OCID you can poll with Get Work Request to confirm removal.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Connect Harness OCID | string | Required | ocid1.connectharness.oc1..aaaa… |
Returns: tool_result, id, work_request_id, success, error
OCI Streaming: Get Connect Harness
oracle/streaming/connect_harness_get · Action
Read a connect harness by OCID — its lifecycle state, name and compartment.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Connect Harness OCID | string | Required | ocid1.connectharness.oc1..aaaa… |
Returns: tool_result, connect_harness, id, lifecycle_state, success, error
OCI Streaming: List Connect Harnesses
oracle/streaming/connect_harness_list · Action
List the connect harnesses in a compartment, optionally filtered by OCID, exact name or lifecycle state. Walks pagination up to a safe cap.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | Required | ocid1.compartment.oc1..aaaa… (use the tenancy OCID for the root) |
| Connect Harness OCID Filter | string | Only the harness with this exact OCID (optional) | |
| Name Filter | string | Only harnesses with this exact name (optional) | |
| Lifecycle State | string | Only harnesses in this state (optional) — choices: Creating, Active, Updating, Deleting, Deleted, Failed | |
| Page Size | string | Items per page, 1–50 (optional) |
Returns: tool_result, connect_harnesses, count, truncated, success, error
OCI Streaming: Update Connect Harness
oracle/streaming/connect_harness_update · Action
Update a connect harness by OCID — change its freeform and defined tags. Only the fields you supply are changed. Returns the harness in an UPDATING state along with a work request OCID.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Connect Harness OCID | string | Required | ocid1.connectharness.oc1..aaaa… |
| Freeform Tags (JSON) | string | {"env":"prod"} (optional) | |
| Defined Tags (JSON) | string | {"Ops":{"env":"prod"}} (optional) |
Returns: tool_result, connect_harness, id, lifecycle_state, work_request_id, success, error
03Consumer
OCI Streaming: Commit Consumer Offsets
oracle/streaming/consumer_commit · Action
Commit a consumer group's cursor offsets so a later reader resumes where the group left off instead of replaying. Supply the stream OCID and the group cursor (from Create Group Cursor); the committed cursor is returned to reuse. The stream's endpoint is resolved automatically.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream OCID | string | Required | ocid1.stream.oc1..aaaa… |
| Group Cursor | text | Required | The group cursor to commit (from Create Group Cursor) |
Returns: tool_result, cursor, success, error
OCI Streaming: Heartbeat Consumer
oracle/streaming/consumer_heartbeat · Action
Heartbeat a group cursor to keep this consumer's partition lease alive so the group does not rebalance away from it. Returns a refreshed cursor to use on the next heartbeat or consume call. The stream's endpoint is resolved automatically.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream OCID | string | Required | ocid1.stream.oc1..aaaa… |
| Group Cursor | text | Required | A group cursor value from Create Group Cursor |
Returns: tool_result, cursor, success, error
04Cursor
OCI Streaming: Create Cursor
oracle/streaming/cursor_create · Action
Create a partition cursor marking where a consumer starts reading — the oldest retained message (TRIM_HORIZON), only new messages (LATEST), a specific offset (AT_OFFSET / AFTER_OFFSET) or a point in time (AT_TIME). Pass the returned cursor into Consume Messages.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream OCID | string | Required | ocid1.stream.oc1..aaaa… |
| Partition | string | Required | The partition to read from, e.g. 0 |
| Cursor Type | string | Required | Where to start reading — choices: Trim Horizon (oldest retained), Latest (only new messages), At Offset, After Offset, At Time |
| Offset | string | Required for At Offset / After Offset, e.g. 42 | |
| Time (RFC 3339) | string | Required for At Time, e.g. 2026-07-22T09:00:00Z |
Returns: tool_result, cursor, success, error
05Group
OCI Streaming: Create Group Cursor
oracle/streaming/group_cursor_create · Action
Create a consumer-group cursor for a stream. Give it a group name and a starting position (LATEST, TRIM_HORIZON, or AT_TIME with a timestamp); the group shares the stream's partitions and tracks its own committed position. Returns a cursor to feed into Consume Messages. The stream's endpoint is resolved automatically.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream OCID | string | Required | ocid1.stream.oc1..aaaa… |
| Group Name | string | Required | Name of the consumer group, e.g. billing-workers |
| Cursor Type | string | Required | Where to start consuming — choices: Latest — only messages published from now on, Trim horizon — the oldest retained message, At time — from a specific timestamp (set Time) |
| Instance Name | string | Unique id for this instance in the group (optional — a UUID is generated) | |
| Timeout (ms) | string | Inactivity before partition reservations are released (optional) | |
| Commit On Get | boolean | Auto-commit each read (default true; set false to commit manually) | |
| Time (AT_TIME only) | string | RFC3339, e.g. 2026-08-01T02:00:00Z — required when Cursor Type is AT_TIME |
Returns: tool_result, cursor, success, error
OCI Streaming: Get Consumer Group
oracle/streaming/group_get · Action
Read the current state of a consumer group on a stream — its partition reservations, the instance holding each partition, and the latest committed offset. The stream's endpoint is resolved automatically.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream OCID | string | Required | ocid1.stream.oc1..aaaa… |
| Consumer Group | string | Required | The name of the consumer group |
Returns: tool_result, group, success, error
OCI Streaming: Reset Consumer Group
oracle/streaming/group_update · Action
Forcefully move a consumer group to a new position in a stream, resetting every consumer at once — to the latest messages, the oldest retained message, or a specific time. The stream's endpoint is resolved automatically.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream OCID | string | Required | ocid1.stream.oc1..aaaa… |
| Consumer Group | string | Required | The name of the consumer group to reset |
| Reset To | string | Required | choices: Latest (only new messages), Trim Horizon (oldest retained message), At Time (a specific timestamp) |
| Time | string | RFC3339, e.g. 2026-12-31T00:00:00Z (required when Reset To is At Time) |
Returns: tool_result, success, error
06Message
OCI Streaming: Consume Messages
oracle/streaming/message_get · Action
Read a batch of messages from a stream using a cursor (from Create Cursor or Create Group Cursor). Returns the messages plus a next cursor to feed into the following call to page forward. The stream's endpoint is resolved automatically.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream OCID | string | Required | ocid1.stream.oc1..aaaa… |
| Cursor | text | Required | A cursor value from Create Cursor / Create Group Cursor |
| Limit | string | Max messages to return, 1–10000 (default service max) |
Returns: tool_result, messages, count, next_cursor, success, error
OCI Streaming: Publish Message
oracle/streaming/message_put · Action
Publish a message to a stream. Supply the stream OCID and a value; an optional key pins related messages to the same partition to preserve their order. The stream's endpoint is resolved automatically.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream OCID | string | Required | ocid1.stream.oc1..aaaa… |
| Message Value | text | Required | The message body (any UTF-8 text or JSON) |
| Message Key | string | Optional partition key — same key ⇒ same partition (ordered) |
Returns: tool_result, partition, offset, timestamp, failures, results, success, error
07Stream
OCI Streaming: Move Stream to Compartment
oracle/streaming/stream_change_compartment · Action
Move a stream into a different compartment. Supply the stream OCID and the destination compartment OCID. The move is asynchronous and returns a work request OCID to poll.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream OCID | string | Required | ocid1.stream.oc1..aaaa… |
| Destination Compartment OCID | string | Required | ocid1.compartment.oc1..aaaa… (where the stream should move to) |
Returns: tool_result, id, destination_compartment_id, work_request_id, success, error
OCI Streaming: Create Stream
oracle/streaming/stream_create · Action
Create a Kafka-compatible stream. Give it a name and a partition count, and either a compartment (uses that compartment's default stream pool) or a specific stream pool. Returns the stream in a CREATING state — poll Get Stream until ACTIVE.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (uses its default stream pool) | |
| Stream Name | string | Required | e.g. orders-events |
| Partitions | string | Required | Number of partitions, e.g. 1 |
| Retention (hours) | string | How long to keep messages, 24–168 (default 24) | |
| Stream Pool OCID | string | ocid1.streampool.oc1..aaaa… (instead of a compartment) | |
| Freeform Tags (JSON) | string | {"env":"prod"} (optional) | |
| Defined Tags (JSON) | string | {"Ops":{"env":"prod"}} (optional) |
Returns: tool_result, stream, id, lifecycle_state, messages_endpoint, work_request_id, success, error
OCI Streaming: Delete Stream
oracle/streaming/stream_delete · Action
Delete a stream by OCID. The stream moves to a DELETING state and its messages are discarded — this cannot be undone. Returns a work request OCID to track the asynchronous removal.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream OCID | string | Required | ocid1.stream.oc1..aaaa… |
Returns: tool_result, id, work_request_id, success, error
OCI Streaming: Get Stream
oracle/streaming/stream_get · Action
Read a stream by OCID — its lifecycle state, partitions, retention and messages endpoint.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream OCID | string | Required | ocid1.stream.oc1..aaaa… |
Returns: tool_result, stream, id, lifecycle_state, messages_endpoint, success, error
OCI Streaming: List Streams
oracle/streaming/stream_list · Action
List the streams in a compartment, optionally narrowed to a stream pool or filtered by name, OCID or lifecycle state. Walks pagination up to a safe cap.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | Required | ocid1.compartment.oc1..aaaa… (use the tenancy OCID for the root) |
| Stream Pool OCID | string | ocid1.streampool.oc1..aaaa… — limit to one pool (optional) | |
| Name Filter | string | Only streams with this exact name (optional) | |
| Stream OCID Filter | string | Only the stream with this exact OCID (optional) | |
| Lifecycle State | string | Only streams in this state (optional) — choices: Creating, Active, Updating, Deleting, Deleted, Failed | |
| Page Size | string | Items per page, 1–50 (default 10) |
Returns: tool_result, streams, count, truncated, success, error
OCI Streaming: Move Stream Pool
oracle/streaming/stream_pool_change_compartment · Action
Move a stream pool into a different compartment. Give it the stream pool OCID and the destination compartment OCID. Asynchronous — returns a work request OCID you can poll.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream Pool OCID | string | Required | ocid1.streampool.oc1..aaaa… |
| Destination Compartment OCID | string | Required | ocid1.compartment.oc1..aaaa… (where to move the pool) |
Returns: tool_result, id, destination_compartment_id, work_request_id, success, error
OCI Streaming: Create Stream Pool
oracle/streaming/stream_pool_create · Action
Create a stream pool in a compartment to hold your streams. Give it a name and optional tags; Kafka settings, private endpoint and encryption are left at their defaults. Returns the pool in a CREATING state — poll Get Stream Pool until ACTIVE.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | Required | ocid1.compartment.oc1..aaaa… (use the tenancy OCID for the root) |
| Stream Pool Name | string | Required | e.g. orders-pool |
| Freeform Tags (JSON) | string | {"env":"prod"} (optional) | |
| Defined Tags (JSON) | string | {"Ops":{"env":"prod"}} (optional) |
Returns: tool_result, stream_pool, id, lifecycle_state, work_request_id, success, error
OCI Streaming: Delete Stream Pool
oracle/streaming/stream_pool_delete · Action
Delete a stream pool by OCID. Asynchronous — the pool moves to DELETING and a work request OCID is returned to track the teardown.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream Pool OCID | string | Required | ocid1.streampool.oc1..aaaa… |
Returns: tool_result, id, work_request_id, success, error
OCI Streaming: Get Stream Pool
oracle/streaming/stream_pool_get · Action
Read a stream pool by OCID — its lifecycle state, compartment, privacy and endpoint FQDN.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream Pool OCID | string | Required | ocid1.streampool.oc1..aaaa… |
Returns: tool_result, stream_pool, id, lifecycle_state, success, error
OCI Streaming: List Stream Pools
oracle/streaming/stream_pool_list · Action
List the stream pools in a compartment, optionally filtered by OCID, exact name, or lifecycle state. Walks pagination up to a safe cap.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | Required | ocid1.compartment.oc1..aaaa… (use the tenancy OCID for the root) |
| Stream Pool OCID Filter | string | Only the pool with this exact OCID (optional) | |
| Name Filter | string | Only pools with this exact name (optional) | |
| Lifecycle State | string | Only pools in this state (optional) — choices: Creating, Active, Updating, Deleting, Deleted, Failed | |
| Page Size | string | Items per page, 1–50 (default 10) |
Returns: tool_result, stream_pools, count, truncated, success, error
OCI Streaming: Update Stream Pool
oracle/streaming/stream_pool_update · Action
Update a stream pool by OCID — change its name and/or its freeform and defined tags. Only the fields you supply are changed. Returns the pool in an UPDATING state along with a work request OCID.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream Pool OCID | string | Required | ocid1.streampool.oc1..aaaa… |
| Stream Pool Name | string | New name (leave blank to keep the current name) | |
| Freeform Tags (JSON) | string | {"env":"prod"} (optional) | |
| Defined Tags (JSON) | string | {"Ops":{"env":"prod"}} (optional) |
Returns: tool_result, stream_pool, id, lifecycle_state, work_request_id, success, error
OCI Streaming: Update Stream
oracle/streaming/stream_update · Action
Update a stream in place. Move it to a different stream pool and/or replace its freeform or defined tags — only the fields you provide are changed. Returns the stream in an UPDATING state; poll Get Stream until ACTIVE.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Stream OCID | string | Required | ocid1.stream.oc1..aaaa… |
| Move to Stream Pool OCID | string | ocid1.streampool.oc1..aaaa… (optional — moves the stream) | |
| Freeform Tags (JSON) | string | {"env":"prod"} — replaces all freeform tags (optional) | |
| Defined Tags (JSON) | string | {"Ops":{"env":"prod"}} — replaces all defined tags (optional) |
Returns: tool_result, stream, id, lifecycle_state, work_request_id, success, error
08Work
OCI Streaming: List Work Request Errors
oracle/streaming/work_request_errors_list · Action
List the errors recorded against a Streaming work request — each error's code, message and timestamp.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Work Request OCID | string | Required | ocid1.streamingworkrequest.oc1..aaaa… |
Returns: tool_result, errors, count, truncated, success, error
OCI Streaming: Get Work Request
oracle/streaming/work_request_get · Action
Read a Streaming work request by OCID — its operation type, status and percent-complete — to track a long-running create, update, delete or move.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Work Request OCID | string | Required | ocid1.streamingworkrequest.oc1..aaaa… |
Returns: tool_result, work_request, status, percent_complete, operation_type, success, error
OCI Streaming: List Work Requests
oracle/streaming/work_request_list · Action
List the asynchronous work requests in a compartment — each tracks a long-running Streaming operation with its operation type, status and percent complete. Walks pagination up to a safe cap.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | Required | ocid1.compartment.oc1..aaaa… (use the tenancy OCID for the root) |
| Page Size | string | Items per page, 1–50 (default 10) |
Returns: tool_result, work_requests, count, truncated, success, error
OCI Streaming: List Work Request Logs
oracle/streaming/work_request_logs_list · Action
List the log entries for a Streaming work request — the progress messages recorded as an asynchronous stream, stream pool or connect harness operation runs. Walks pagination up to a safe cap.
| Field | Type | Details | |
|---|---|---|---|
| Compartment OCID | string | ocid1.compartment.oc1..aaaa… (scopes the picker) | |
| Work Request OCID | string | Required | ocid1.streamingworkrequest.oc1..aaaa… |
Returns: tool_result, logs, count, truncated, success, error
09Notes & Limitations
Behaviours and constraints worth knowing before you build with these nodes.
- Publishing, consuming, cursors and consumer groups each reach a stream on its own messages endpoint, which Flomation resolves automatically from the stream's OCID, so a stream that is still provisioning is rejected until it reaches
ACTIVErather than being queued. - Creating, updating, deleting and moving a stream all run asynchronously and return a work request OCID, so let one operation finish — polling the matching Get action until the stream settles — before starting another change against the same stream.
- Creating a stream takes exactly one of a compartment OCID or a stream pool OCID: a compartment places the stream in that compartment's default stream pool, while a pool OCID puts it in that specific pool.
- A stream's retention period is fixed when it is created and cannot be changed afterwards, so set the Retention (hours) value — any whole number of hours from 24 to 168 — before you create the stream.
- A published message value is stored as exact bytes, so any leading or trailing whitespace or newlines you include are returned unchanged when the message is consumed.