Skip to main content

OCI Streaming — Kafka-Compatible Streaming Service

Calculating read time…

You've built a Kafka-powered streaming pipeline. Your producers publish events. Your consumers process them in real time. Everything works beautifully — but now you're on Oracle Cloud Infrastructure (OCI), and the last thing you want is to maintain a self-managed Kafka cluster: ZooKeeper configuration, broker upgrades, replication tuning, disk management, scaling scripts, 3 AM alerts.

Enter OCI Streaming — Oracle's fully managed, serverless streaming service that is completely compatible with the Apache Kafka API. Your existing Kafka producers and consumers connect to OCI Streaming with a single configuration change — no code modifications required.

This post is the most detailed technical guide to OCI Streaming available. We will map every Kafka concept to its OCI Streaming equivalent, show you exactly how to configure authentication, walk through real production configurations, and reveal the subtle differences you must know before going live.





💡 OCI Streaming — The Key Numbers (2026)

📨 Kafka API compatible — uses standard Kafka protocol (port 9092)
🔀 1–25 partitions per stream (increase via service request)
📦 1 MB max message size
⏱️ 7 days maximum message retention (168 hours)
🚀 1 MB/s throughput per partition (write) · 2 MB/s (read)
🌍 Available in all major OCI regions: us-ashburn-1, eu-frankfurt-1, ap-tokyo-1, ap-mumbai-1, etc.
💰 Pricing: per partition-hour + data in/out (egress)
🔌 No ZooKeeper, no broker management — fully managed by Oracle

The bootstrap server format: cell-1.streaming.{region}.oci.oraclecloud.com:9092

📨 Section 1: Apache Kafka — A 60-Second Recap

Before we map OCI Streaming to Kafka, let's be crystal clear about what Kafka is. This ensures every comparison that follows makes perfect sense.

💡 The Highway System Analogy for Kafka

Kafka is like a massive, multi-lane highway interchange. Producers are vehicles entering the highway, carrying cargo (messages). Topics are specific roads — the "I-90 Eastbound" of your data. Partitions are individual lanes on that road — messages are spread across lanes for parallelism. Consumers are vehicles exiting at the other end, picking up cargo. Consumer Groups are fleets of vehicles working together, each handling different lanes.

Kafka is the infrastructure — durable, ordered, replayable, massively scalable.

📨 Apache Kafka — Core Architecture at a Glance

🖥️ Producers
Write messages
to topics
→
📨 Kafka Cluster
Topic: "orders"
P0
P1
P2
3 partitions | 7-day retention
→
👥 Consumer Group A
Reads P0+P1 (analytics)
👥 Consumer Group B
Reads P0+P1+P2 (notifications)

↑ Multiple consumer groups can read the same topic independently. Each starts from its own offset. This is why Kafka is called a "distributed commit log" — it remembers everything for every consumer.


☁️ Section 2: What Is OCI Streaming — Oracle's Managed Kafka

OCI Streaming is Oracle Cloud's fully managed, Apache Kafka-compatible real-time message streaming service. Oracle manages the entire infrastructure: brokers, replication, upgrades, scaling, backups, and availability. You consume the service — Oracle runs it.

🔴 The Hotel vs Airbnb Analogy

Self-managed Kafka is like buying and running your own apartment building. You own every machine. You're responsible for: hardware failures, OS patches, Kafka upgrades, ZooKeeper configuration, disk expansion, network tuning, broker rebalancing. Full control — full responsibility. Great if you have a dedicated Kafka team.

OCI Streaming is like staying at a five-star hotel. You get all the capabilities — your own room (stream), key (partition), storage (retention). Oracle handles all the plumbing. You just use the service and pay per night (per partition-hour). Perfect for teams that want Kafka without the operational overhead.
❌ Self-Managed Kafka — You Must Handle
🖥️ Provision and configure Kafka brokers
🐘 ZooKeeper / KRaft cluster management
📀 Disk management and expansion
🔁 Replication factor configuration
🔧 Kafka version upgrades
📊 JVM tuning, OS tuning, GC pauses
🔔 3 AM broker failure alerts
💰 Paying for idle broker capacity
✅ OCI Streaming — Oracle Handles All This
🏗️ Infrastructure provisioning
🔁 Automatic replication and HA
📀 Storage auto-scaling
🔧 Version upgrades (zero downtime)
📊 Metrics and monitoring built-in
🔐 Security (TLS, IAM integrated)
🌍 Multi-AD high availability
You only manage: Streams, Partitions, Messages ✅

🔄 Section 3: The Complete Kafka → OCI Streaming Concept Mapping

This is the heart of this post. Every Kafka concept has a direct equivalent in OCI Streaming. Once you understand these mappings, you understand OCI Streaming completely.

🔄 Apache Kafka ≡ OCI Streaming — Complete Concept Equivalence

📦
Kafka Cluster
A group of brokers running together. Namespace for all your topics.
≡
🏊
Stream Pool
Container for streams. Defines shared settings (retention, Kafka settings, encryption).
📰
Kafka Topic
Named, ordered, durable log. Messages are appended and retained for configured TTL.
≡
🌊
Stream
Named message channel inside a Stream Pool. Has configurable partitions and retention.
🔀
Kafka Partition
Ordered sub-log within a topic. Unit of parallelism. Each has its own offset sequence.
≡
🔀
Partition (same name!)
Identical concept. 1–25 per stream. Throughput: 1 MB/s write, 2 MB/s read per partition.
📩
Kafka Record
Key + Value + Timestamp + Headers. Max 1 MB default. Immutable once written.
≡
📩
Message
Key + Value + Timestamp. Max 1 MB. Written via Kafka protocol or OCI REST API.
👥
Consumer Group
Named group of consumers. Each partition assigned to exactly one consumer in the group.
≡
👥
Stream Group
OCI term for Consumer Group. Via Kafka client: group.id config. Via REST API: explicit group cursor.
📍
Kafka Offset
Monotonically increasing integer per partition. Uniquely identifies each record's position.
≡
📍
Offset / Cursor
Via Kafka API: same integer offset. Via REST API: cursor string (opaque token) for sequential reads.
⏱️
retention.ms
How long messages are kept. Configurable per topic. Default 7 days. Can be unlimited.
≡
⏱️
Retention Period
Set per Stream: 24 hours (default) up to 7 days (168 hours). Cannot exceed 7 days currently.
🖥️
Kafka Broker
Individual server in the Kafka cluster. Stores partitions, handles replication, serves clients.
≡
🌐
Oracle-Managed Nodes
Fully managed by Oracle. You never interact with brokers. The bootstrap server abstracts them all.

🏊 Section 4: Stream Pools — The OCI Kafka Cluster

The Stream Pool is the most important OCI-specific concept to understand. It is what OCI Streaming uses in place of a Kafka "cluster" or "namespace."

🏊 Stream Pool — Structure and Configuration

📍 Location: Lives within an OCI Compartment and Region. One stream pool = one Kafka cluster namespace.
🔐 Kafka Settings: Each Stream Pool has a Kafka Settings section that gives you: the bootstrap server, the SASL connection string, and the auth token instructions. This is where you get everything you need to connect a Kafka client.
🔒 Encryption: All data encrypted at rest (using OCI Vault or Oracle-managed keys). All data encrypted in transit (TLS 1.2+). You never handle certificates manually.
🌍 Availability: Data is automatically replicated across multiple Availability Domains in the OCI region. No replication factor to configure — Oracle handles it.
🔑 Private Endpoint: Stream Pools can optionally have a private endpoint (only accessible from within a VCN). This prevents any public internet exposure — essential for enterprise security requirements.

🏊 Inside a Stream Pool — Contains Multiple Streams

🏊 Stream Pool: "production-events-pool"
🌊 Stream:
"orders"
5 partitions
7-day retention
🌊 Stream:
"payments"
3 partitions
7-day retention
🌊 Stream:
"user-events"
10 partitions
24h retention
🌊 Stream:
"notifications"
1 partition
1-day retention
All streams share the pool's private endpoint, encryption key, and IAM policies. Bootstrap server: cell-1.streaming.ap-mumbai-1.oci.oraclecloud.com:9092

🔌 Section 5: Kafka API Compatibility — Zero Code Changes

This is the most important practical feature of OCI Streaming. OCI Streaming implements the Kafka wire protocol on port 9092. Any application using a standard Kafka client library can connect to OCI Streaming by simply changing the bootstrap server address and authentication configuration. No code changes. No API changes. Same Kafka client.

✅ Kafka Operations Supported by OCI Streaming

✅ Producer API (produce messages)
✅ Consumer API (subscribe & poll)
✅ Consumer Groups & rebalancing
✅ Offset management (commit, seek)
✅ Message keys for partitioning
✅ Message headers
✅ SASL/PLAIN authentication
✅ TLS encryption (port 9092)
✅ List topics (metadata API)
✅ Kafka Connect (source/sink)
✅ Kafka Streams applications
✅ Mirror Maker 2 (replication)

❌ Not Currently Supported in OCI Streaming

❌ Kafka Transactions (exactly-once)
❌ Log compaction (compact cleanup)
❌ Kafka Admin API (create topics via Kafka protocol — use OCI API/Console)
❌ Schema Registry (use OCI Schema Registry separately)
❌ Retention beyond 7 days
❌ Custom Kafka configurations (broker-level settings)

🔐 Section 6: Authentication — OCI IAM Meets Kafka SASL

Authentication is where OCI Streaming differs most visibly from a self-managed Kafka cluster. OCI Streaming uses IAM (Identity and Access Management) for access control, translated into Kafka's SASL/PLAIN mechanism. Understanding this translation is essential for a successful connection.

💡 The OCI Auth Token Analogy

In self-managed Kafka, you create a Kafka-specific username and password. In OCI Streaming, you use your OCI identity (user or service principal) + an Auth Token (a secret token you generate in the OCI Console under your profile → Auth Tokens).

The Auth Token is essentially a permanent API password tied to your OCI user account. It's passed as the Kafka SASL password. The username encodes your tenancy and identity. This way, OCI IAM policies control who can produce/consume from which streams.

🔐 OCI Streaming Authentication — Three Identity Types

👤 Type 1: User Principal (most common for development)
SASL Username = {tenancy_name}/{username}
Example: mycompany/john.doe@example.com
SASL Password = Auth Token generated in OCI Console → Profile → Auth Tokens
⚠️ Not suitable for production automation. Auth tokens are user-scoped and expire.
🤖 Type 2: OCI Service (Instance Principal / Resource Principal) — Production Best Practice
Compute instances, OCI Functions, Container Instances automatically get identity credentials.
SASL Username = {tenancy_ocid}/<instance-principal>
SASL Password = Dynamically generated token (retrieved via OCI SDK, auto-refreshed)
✅ No hardcoded credentials. Tokens auto-rotate. The secure production approach.
🔑 Type 3: OCI CLI / SDK Auth (for admin operations via REST API)
For creating streams, managing pools via OCI REST API or CLI — uses OCI API key (private key + key fingerprint). This is separate from the Kafka client connection. Admin = API key auth. Data plane = SASL auth.

🔐 Kafka JAAS Configuration for OCI Streaming

# OCI Streaming Kafka Client Configuration (kafka.properties) # This is ALL you change to point your Kafka app to OCI Streaming bootstrap.servers=cell-1.streaming.ap-mumbai-1.oci.oraclecloud.com:9092 security.protocol=SASL_SSL sasl.mechanism=PLAIN sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required \ username="mytenancy/john.doe@example.com" \ password="YOUR_OCI_AUTH_TOKEN_HERE"; # Standard Kafka settings — unchanged from regular Kafka key.serializer=org.apache.kafka.common.serialization.StringSerializer value.serializer=org.apache.kafka.common.serialization.StringSerializer key.deserializer=org.apache.kafka.common.serialization.StringDeserializer value.deserializer=org.apache.kafka.common.serialization.StringDeserializer

↑ That is the entire configuration change to migrate from self-managed Kafka to OCI Streaming. Your producer and consumer application code doesn't change at all. Just these config values.


🔀 Section 7: Partitions, Retention and Throughput — OCI Streaming Limits

Understanding OCI Streaming's limits is essential for designing a production system. Most limits are generous for typical workloads, but a few require planning.

🔀 Partitions — The Unit of Parallelism

🔀 OCI Streaming Partition Model — Identical to Kafka

Partition 0
Msg offset 0
Msg offset 1
Msg offset 2
→ offset 3 (next)
Partition 1
Msg offset 0
Msg offset 1
→ offset 2 (next)
Partition 2
Msg offset 0
Msg offset 1
→ offset 2 (next)
OCI Streaming Partition Rules (same as Kafka):
✅ Messages with the same key always go to the same partition (consistent ordering per key)
✅ Messages without a key are distributed round-robin across partitions
✅ Each partition is consumed by exactly one consumer within a consumer group
✅ Multiple consumer groups can consume the same partition independently
⚠️ Max partitions per stream: 25 (default). Can be increased via service limit request.

⏱️ Message Retention — OCI's 7-Day Maximum

⏱️ Retention Period — Message Lifecycle in OCI Streaming

Messages are retained for the configured duration and then automatically deleted:
Default:
24 hours
Recommended:
72 hours (3 days)
Maximum:
168 hours (7 days)
⚠️ Important: Unlike self-managed Kafka (where you can keep messages indefinitely or use log compaction), OCI Streaming has a hard 7-day (168-hour) maximum retention. If your use case requires longer retention (e.g., event sourcing, long-term audit), you must export messages to OCI Object Storage using the Service Connector Hub. The Service Connector Hub can automatically archive all messages to Object Storage as they arrive.

🚀 Throughput Limits Per Partition

📊 Metric 📤 OCI Streaming (per partition) 📨 Kafka Default (per partition)
Write throughput 1 MB/s Configurable (hardware-limited)
Read throughput 2 MB/s Configurable (hardware-limited)
Max message size 1 MB 1 MB default (configurable)
Max partitions/stream 25 (default limit) No limit (broker capacity)
Total throughput per stream 25 MB/s write, 50 MB/s read
(with 25 partitions)
Scale horizontally
Retention Max 7 days (168 hours) Unlimited (log compaction supported)
✅ How to Scale Beyond Throughput Limits:

Need more than 1 MB/s write per partition? Add more partitions. Need 100 MB/s total write throughput? Use 100 partitions (requires service limit increase request).

Need more than 7-day retention? Use Service Connector Hub to continuously archive messages to OCI Object Storage. Your consumers can still read the live 7-day window, and long-term data lives in Object Storage (essentially unlimited, much cheaper).

🔌 Section 8: Two Ways to Access OCI Streaming — Kafka vs REST API

OCI Streaming offers two completely different API surfaces. Understanding both is important for choosing the right integration pattern.

📨 Method 1: Kafka Protocol API (Port 9092)

Use any standard Kafka client library (Java, Python, Go, .NET, etc.). Uses the native Kafka binary wire protocol. Full feature set of supported Kafka operations.

✅ Best for: existing Kafka applications
✅ Zero code changes to migrate
✅ Consumer groups, offsets, rebalancing
✅ All Kafka client features
⚠️ Requires: SASL_SSL config + auth token

🌐 Method 2: OCI REST API (HTTPS)

Use OCI's HTTP REST API (streaming.{region}.oci.oraclecloud.com). Authenticated via OCI API Key Signature (standard OCI API auth).

✅ Best for: serverless, functions, quick integrations
✅ No Kafka client library needed
✅ Works from any HTTP client (curl, Postman)
✅ OCI SDK support (Java, Python, Go, etc.)
⚠️ Uses cursor-based reading (not Kafka offset API)
⚠️ No consumer groups via REST (manual cursor mgmt)

💻 Section 9: Complete Code Examples — Producer, Consumer, REST API

📌 What This Code Does (Read Before The Code!)

This is a Python Kafka producer that connects to OCI Streaming using the standard kafka-python library — the exact same library you would use with a self-managed Kafka cluster. The ONLY difference from a regular Kafka producer is the connection configuration: the bootstrap_servers points to OCI's endpoint, and security_protocol, sasl_mechanism, sasl_plain_username, and sasl_plain_password are set to OCI IAM credentials (tenancy/username + auth token). The actual producer code — creating a KafkaProducer, sending messages with keys, handling delivery callbacks — is 100% identical to what you would write for any Kafka cluster. This demonstrates OCI Streaming's Kafka API compatibility in the most concrete way possible.

# ✅ Python Kafka Producer → OCI Streaming
# Uses standard kafka-python library. ZERO code changes vs regular Kafka.
# pip install kafka-python

from kafka import KafkaProducer
import json, time

# ── OCI Streaming Connection Configuration ─────────────────────
# This is the ONLY difference from a regular Kafka producer.
OCI_BOOTSTRAP   = "cell-1.streaming.ap-mumbai-1.oci.oraclecloud.com:9092"
OCI_TENANCY     = "mytenancy"
OCI_USERNAME    = "john.doe@example.com"
OCI_AUTH_TOKEN  = "YOUR_OCI_AUTH_TOKEN_HERE"  # from OCI Console → Profile → Auth Tokens
OCI_STREAM_NAME = "orders"                   # OCI Stream = Kafka Topic

# ── Create Kafka Producer with OCI authentication ──────────────
producer = KafkaProducer(
    # OCI Streaming endpoint (same port 9092 as Kafka!)
    bootstrap_servers     = [OCI_BOOTSTRAP],

    # OCI uses SASL_SSL — encrypts + authenticates
    security_protocol     = "SASL_SSL",
    sasl_mechanism        = "PLAIN",
    sasl_plain_username   = f"{OCI_TENANCY}/{OCI_USERNAME}",
    sasl_plain_password   = OCI_AUTH_TOKEN,

    # Standard Kafka producer settings — unchanged!
    value_serializer      = lambda v: json.dumps(v).encode("utf-8"),
    key_serializer        = lambda k: k.encode("utf-8"),

    # Performance tuning (same as regular Kafka)
    acks                  = "all",      # wait for all replicas to ACK
    retries               = 3,
    batch_size            = 16384,
    linger_ms             = 10,
)

# ── Produce messages (identical to regular Kafka code!) ─────────
orders = [
    {"order_id": "ORD-001", "user_id": "user-42", "amount": 1500.00, "status": "placed"},
    {"order_id": "ORD-002", "user_id": "user-99", "amount": 899.50,  "status": "placed"},
]

for order in orders:
    # key = order_id ensures all events for one order go to same partition
    future = producer.send(
        topic  = OCI_STREAM_NAME,    # OCI Stream name = Kafka topic name
        key    = order["order_id"],  # routing key for partition assignment
        value  = order              # auto-serialized to JSON bytes
    )
    record_metadata = future.get(timeout=10)  # wait for ACK
    print(
        f"✅ Published: {order['order_id']} → "
        f"partition={record_metadata.partition}, offset={record_metadata.offset}"
    )

producer.flush()  # ensure all buffered messages are sent
producer.close()
📌 What This Code Does (Read Before The Code!)

This is a Python Kafka consumer that reads from OCI Streaming as part of a consumer group. Just like with the producer, the only changes from a regular Kafka consumer are in the connection configuration. The consumer subscribes to the "orders" stream (= Kafka topic), joins consumer group "order-processor-group", and polls for messages in an infinite loop. OCI Streaming manages the partition assignment and offset tracking automatically — exactly as Kafka does. When you restart this consumer, it picks up from where it left off (committed offset). If you run two instances of this consumer with the same group_id, OCI Streaming (via Kafka consumer group protocol) automatically distributes partitions between them — one consumer per partition. This is the same behavior as Apache Kafka's consumer group rebalancing.

# ✅ Python Kafka Consumer ← OCI Streaming
# Consumer group with auto-commit and partition rebalancing.
# 100% standard Kafka consumer code with OCI connection config.

from kafka import KafkaConsumer
import json

OCI_BOOTSTRAP   = "cell-1.streaming.ap-mumbai-1.oci.oraclecloud.com:9092"
OCI_TENANCY     = "mytenancy"
OCI_USERNAME    = "john.doe@example.com"
OCI_AUTH_TOKEN  = "YOUR_OCI_AUTH_TOKEN_HERE"

# ── Create Kafka Consumer connected to OCI Streaming ──────────
consumer = KafkaConsumer(
    "orders",                     # OCI Stream name (= Kafka topic)

    # OCI Streaming connection + auth (the only OCI-specific part)
    bootstrap_servers     = [OCI_BOOTSTRAP],
    security_protocol     = "SASL_SSL",
    sasl_mechanism        = "PLAIN",
    sasl_plain_username   = f"{OCI_TENANCY}/{OCI_USERNAME}",
    sasl_plain_password   = OCI_AUTH_TOKEN,

    # Consumer group — identical to Kafka consumer groups!
    # All consumers with the same group_id share partitions (load balanced)
    # Each partition goes to exactly one consumer in the group.
    group_id              = "order-processor-group",

    # Start from earliest message if no committed offset exists yet
    auto_offset_reset     = "earliest",

    # Auto-commit offsets every 5 seconds
    enable_auto_commit    = True,
    auto_commit_interval_ms = 5000,

    # Deserialize messages from JSON bytes
    value_deserializer    = lambda m: json.loads(m.decode("utf-8")),
    key_deserializer      = lambda k: k.decode("utf-8") if k else None,

    # Session timeout for consumer group coordination
    session_timeout_ms    = 30000,
    heartbeat_interval_ms = 10000,
)

print("🟢 Consumer started. Waiting for messages from OCI Streaming...")

# ── Poll for messages (infinite loop — same as regular Kafka) ──
for message in consumer:
    order = message.value

    print(
        f"📨 Received from partition={message.partition}, offset={message.offset}\n"
        f"   Order: {order['order_id']} | User: {order['user_id']} | ₹{order['amount']}"
    )

    # Your business logic here — no different from any Kafka consumer
    # e.g., process_order(order), call_payment_service(order), etc.

    # Manual commit (if enable_auto_commit=False):
    # consumer.commit()  ← commit after successful processing
📌 What This Code Does (Read Before The Code!)

This shows OCI Streaming's native REST API — the alternative to the Kafka protocol approach. Instead of using a Kafka client library, you use OCI's HTTP API with OCI SDK authentication (API key signature). This is ideal for: OCI Functions (serverless), quick scripts, languages without a good Kafka client, or situations where you want to avoid the Kafka client dependency. The REST API has two key differences from Kafka: (1) Reading uses a "cursor" concept instead of offsets (you first create a cursor specifying where to start, then use it to get messages), and (2) Writing uses base64-encoded message values. The OCI Python SDK handles all the authentication signing automatically — you just provide the config file path. This exact pattern is used by OCI Functions to consume streaming events without needing a full Kafka client.

# ✅ OCI Streaming REST API — Using OCI Python SDK
# Alternative to Kafka protocol. No Kafka client library needed!
# pip install oci

import oci, base64, json

# ── OCI Configuration ──────────────────────────────────────────
config = oci.config.from_file("~/.oci/config")  # OCI CLI config with API key

# Stream OCID from OCI Console → Streaming → your stream
STREAM_OCID    = "ocid1.stream.oc1.ap-mumbai-1.amaaaaa..."
STREAM_ENDPOINT= "https://cell-1.streaming.ap-mumbai-1.oci.oraclecloud.com"

# Create streaming client
streaming_client = oci.streaming.StreamClient(config, service_endpoint=STREAM_ENDPOINT)
admin_client     = oci.streaming.StreamAdminClient(config)

# ════════════════════════════════════════════════════════════════
# PRODUCE: Publish messages via REST API
def publish_messages(messages: list):
    # OCI REST API requires base64-encoded values
    put_messages_details = oci.streaming.models.PutMessagesDetails(
        messages=[
            oci.streaming.models.PutMessagesDetailsEntry(
                key   = base64.b64encode(msg["key"].encode()).decode(),
                value = base64.b64encode(
                    json.dumps(msg["value"]).encode()
                ).decode()
            )
            for msg in messages
        ]
    )

    response = streaming_client.put_messages(
        stream_id           = STREAM_OCID,
        put_messages_details= put_messages_details
    )

    for entry in response.data.entries:
        if entry.error:
            print(f"❌ Failed: {entry.error_message}")
        else:
            print(f"✅ Published → partition={entry.partition}, offset={entry.offset}")


# ════════════════════════════════════════════════════════════════
# CONSUME: Read messages via REST API using cursor
# Note: REST API uses CURSOR (not Kafka offset) for position tracking
def consume_messages(partition: int = 0, limit: int = 10):

    # Step 1: Create a cursor — define WHERE to start reading
    # AT_OFFSET: start at a specific offset (like Kafka seek)
    # LATEST:    read only new messages (like auto_offset_reset=latest)
    # EARLIEST:  read from beginning  (like auto_offset_reset=earliest)
    cursor_details = oci.streaming.models.CreateCursorDetails(
        partition      = str(partition),
        type           = oci.streaming.models.CreateCursorDetails.TYPE_TRIM_HORIZON,
        # TYPE_TRIM_HORIZON = "earliest" in Kafka terminology
    )

    cursor_response = streaming_client.create_cursor(
        stream_id           = STREAM_OCID,
        create_cursor_details= cursor_details
    )
    cursor = cursor_response.data.value  # opaque cursor string

    # Step 2: Use cursor to read messages (poll loop)
    while True:
        get_response = streaming_client.get_messages(
            stream_id = STREAM_OCID,
            cursor    = cursor,
            limit     = limit
        )

        messages = get_response.data
        if not messages:
            print("📭 No new messages. Waiting...")
            continue

        for msg in messages:
            # Values are base64-encoded — decode them
            value = json.loads(base64.b64decode(msg.value).decode("utf-8"))
            key   = base64.b64decode(msg.key).decode("utf-8") if msg.key else None
            print(f"📨 partition={msg.partition}, offset={msg.offset}, key={key}")
            print(f"   Value: {value}")

        # Move cursor forward for next poll
        cursor = get_response.headers.get("opc-next-cursor")


# Example usage:
publish_messages([
    {"key": "ORD-001", "value": {"order_id": "ORD-001", "amount": 1500}},
    {"key": "ORD-002", "value": {"order_id": "ORD-002", "amount": 899}},
])
consume_messages(partition=0)
📌 What This Code Does (Read Before The Code!)

This shows a Spring Boot application configured for OCI Streaming. Spring Boot's Kafka integration (spring-kafka) uses standard application.yml configuration. In a regular Kafka setup, you'd set bootstrap-servers to your Kafka cluster address. Here you set it to OCI Streaming's endpoint, add SASL security properties, and everything else — @KafkaListener, KafkaTemplate, consumer group, serializers — works exactly as it does with regular Kafka. No Spring Boot code changes needed. This is the configuration pattern used by Java microservices on OCI when migrating from self-managed Kafka to OCI Streaming. The configuration changes are confined entirely to application.yml — all Java @Service, @Component annotations remain unchanged.

# ✅ Spring Boot — application.yml for OCI Streaming
# Java microservices migrate to OCI Streaming with ONLY config changes.
# All @KafkaListener and KafkaTemplate code remains 100% unchanged.

spring:
  kafka:
    bootstrap-servers: cell-1.streaming.ap-mumbai-1.oci.oraclecloud.com:9092

    # OCI Streaming authentication via SASL
    properties:
      security.protocol: SASL_SSL
      sasl.mechanism: PLAIN
      sasl.jaas.config: >
        org.apache.kafka.common.security.plain.PlainLoginModule required
        username="mytenancy/john.doe@example.com"
        password="${OCI_AUTH_TOKEN}";  # from environment variable

    producer:
      key-serializer:   org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
      acks: all
      retries: 3

    consumer:
      group-id: order-service-group
      auto-offset-reset: earliest
      enable-auto-commit: false  # manual ack for reliability
      key-deserializer:   org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
      properties:
        spring.json.trusted.packages: "com.myapp.model.*"

---
# Your Java code — COMPLETELY UNCHANGED from regular Kafka!
# @KafkaListener, KafkaTemplate — identical to self-managed Kafka.

// OrderEventProducer.java
@Service
public class OrderEventProducer {
    @Autowired
    private KafkaTemplate<String, OrderEvent> kafkaTemplate;

    public void publishOrderCreated(Order order) {
        // Sends to OCI Streaming "orders" stream — same as Kafka topic!
        kafkaTemplate.send("orders", order.getId(), new OrderEvent(order));
        log.info("Published order {} to OCI Streaming", order.getId());
    }
}

// OrderEventConsumer.java
@Component
public class OrderEventConsumer {

    @KafkaListener(topics = "orders", groupId = "order-service-group")
    public void consumeOrderEvent(
            OrderEvent event,
            @Header(KafkaHeaders.RECEIVED_PARTITION_ID) int partition,
            @Header(KafkaHeaders.OFFSET) long offset,
            Acknowledgment ack) {

        // Process order — this code is 100% identical to regular Kafka!
        processOrder(event);
        ack.acknowledge();  // manual offset commit after processing
    }
}
✅ Migration From Self-Managed Kafka to OCI Streaming — The 3-Step Process:

Step 1: Create a Stream Pool and Streams in OCI Console (matching your existing Kafka topic names)
Step 2: Generate an Auth Token for your OCI user or Service Principal
Step 3: Update ONLY the connection configuration (bootstrap_servers, sasl config) in your application

Application code: unchanged. Business logic: unchanged. Test suite: unchanged. Deploy. Done. Migration complete. 🎉

📊 Section 10: OCI Streaming vs Self-Managed Kafka — Detailed Comparison

Dimension 🔴 OCI Streaming (Managed) 🟢 Self-Managed Kafka (on OCI/on-prem)
Setup time ✅ Minutes (console/Terraform) Days (broker config, ZooKeeper, networking)
Operational overhead ✅ Zero — Oracle manages everything High — dedicated ops team required
Kafka API compatibility ✅ Full producer/consumer API ✅ Complete Kafka API
Message retention Max 7 days ✅ Unlimited (disk permitting)
Log compaction ❌ Not supported ✅ Supported
Exactly-once (transactions) ❌ Not supported ✅ Kafka Transactions API
High availability ✅ Automatic, multi-AD Manual: configure replication factor, rack awareness
Scaling Add partitions (limit: 25 default) ✅ Add brokers, unlimited partitions
OCI Integration ✅ Native: Functions, Service Connector Hub, GoldenGate Manual integration required
Security ✅ OCI IAM integrated, TLS auto-managed Manual: TLS certs, SASL config, ACLs
Cost model Pay per partition-hour + egress Pay for VMs, storage, networking (potentially cheaper at extreme scale)
Best for Teams on OCI wanting Kafka without ops overhead. Moderate scale. OCI-native integrations. Large scale. Unlimited retention. Transactions. Log compaction. Maximum control.

🔗 Section 11: OCI Streaming Native Integrations

One of the biggest advantages of OCI Streaming over self-managed Kafka is its deep integration with other OCI services through the Service Connector Hub.

🔗 OCI Service Connector Hub — Connect Streaming to Any OCI Service

🌊 OCI Streaming Source
Your stream data flows continuously
⬇️ Service Connector Hub (no-code pipeline)
📦 OCI Object Storage
Archive messages
forever cheaply
⚡ OCI Functions
Run serverless
on each event
🔔 OCI Notifications
SMS / Email
on events
🔍 OCI Logging
Analytics &
audit trails
📊 OCI Monitoring
Custom metrics
from events

↑ Each connector is configured visually in the OCI Console — zero code. This is how OCI Streaming overcomes the 7-day retention limit: stream to Object Storage and keep messages forever.

💡 OCI GoldenGate + OCI Streaming = Enterprise CDC

OCI GoldenGate is Oracle's enterprise Change Data Capture platform. It captures every change from Oracle Database, MySQL, PostgreSQL, and many others and streams them directly into OCI Streaming topics.

This is the OCI equivalent of the Debezium + Kafka pattern: database changes → OCI Streaming → downstream consumers (Elasticsearch, data warehouse, microservices).

Zero code required. GUI-configured pipelines. Enterprise-grade reliability with Oracle support.

🗺️ Section 12: Everything Together — OCI Streaming Production Architecture

☁️ OCI Streaming — Complete Production Architecture

── PRODUCERS (use standard Kafka client libraries) ──
🖥️ Java Spring Boot Apps
🐍 Python Microservices
🔧 OCI GoldenGate (CDC)
⚡ OCI Functions
⬇️ SASL_SSL → port 9092
☁️ OCI STREAMING (Managed Kafka)
🏊 Stream Pool: "production-pool"
🌊 orders (5 partitions, 7d)
🌊 payments (3 partitions, 7d)
🌊 user-events (10 partitions, 24h)
🔐 OCI IAM Auth
🔒 TLS Encryption
🌍 Multi-AD HA
📊 OCI Monitoring
⬇️
── CONSUMERS ──
👥 Consumer Group A
(Order Processing)
👥 Consumer Group B
(Analytics Flink)
👥 Consumer Group C
(Notification Svc)
⬇️ Service Connector Hub
📦 Object Storage
(long-term archive)
⚡ OCI Functions
(serverless processing)
🔔 OCI Notifications
(alerts)
📊 OCI Logging
(audit)

📐 Section 13: OCI Streaming Best Practices

🏊 One Stream Pool Per Environment

Create separate Stream Pools for dev, staging, and production. This provides clean isolation — dev consumers can never accidentally consume prod messages. Stream Pools are also where you configure encryption keys, private endpoints, and Kafka settings. Having per-environment pools makes these configurations independent and auditable.

⏱️ Always Configure Maximum Retention + Service Connector Hub Archive

Set retention to 7 days (168 hours) on all streams — the maximum. This gives you the maximum replay window for debugging, re-processing, and consumer restarts. Simultaneously, set up a Service Connector Hub rule to archive all messages to OCI Object Storage. This gives you effectively unlimited retention history in Object Storage at minimal cost, while keeping 7 days always available in the live stream for immediate consumption.

🔑 Use Instance Principal Auth in Production (Never Hardcode Tokens)

Auth tokens are user-scoped and can expire. For production applications on OCI Compute, OCI Container Engine (OKE), or OCI Functions, use Instance Principal or Resource Principal authentication. These automatically provide rotating credentials without any hardcoded secrets. Create an IAM Dynamic Group for your compute instances and grant it streaming permissions via policy.

📊 Monitor These Three OCI Streaming Metrics

(1) PutMessagesRequestsPerSecond — write throughput per stream (alarm if near 1 MB/s × partitions limit)
(2) GetMessagesRequestsPerSecond — read throughput per consumer group
(3) ConsumerGroupLag — how far behind each consumer group is (alarm if lag grows unexpectedly)
All available in OCI Monitoring with alarm capabilities.

🔑 Message Keys for Ordering Guarantees

Always set a message key for operations that must be ordered for a specific entity. OCI Streaming (like Kafka) guarantees ordering only within a partition. If all events for order "ORD-001" must be processed in order, use "ORD-001" as the key — all events will always go to the same partition and be processed in the order they were written.


🎓 Section 14:  Cheat Sheet

❓ "What is OCI Streaming and how is it related to Kafka?"

OCI Streaming is Oracle Cloud's fully managed, serverless streaming service that is Kafka API-compatible. It implements the Apache Kafka wire protocol on port 9092, so any existing Kafka producer or consumer application can connect to OCI Streaming by changing only the bootstrap server address and authentication configuration — no code changes. Oracle manages all infrastructure: brokers, replication, upgrades, HA. Core concepts map 1:1: Stream Pool = Kafka Cluster, Stream = Topic, Partition = Partition, Stream Group = Consumer Group.

❓ "What are the key differences between OCI Streaming and self-managed Kafka?"

Limitations of OCI Streaming vs full Kafka: (1) Max 7-day retention (no unlimited). (2) No log compaction. (3) No Kafka Transactions (no exactly-once). (4) Max 25 partitions per stream by default. (5) 1 MB/s write per partition throughput limit. (6) Topics created via OCI console/API, not Kafka Admin API. Advantages: Zero operational overhead, automatic HA, native OCI IAM security, Service Connector Hub integration, free TLS management, minutes to set up vs days for self-managed.

❓ "How does authentication work in OCI Streaming?"

OCI Streaming uses SASL/PLAIN over TLS (SASL_SSL). Username = tenancyName/username (or tenancyOCID/userOCID). Password = OCI Auth Token (generated in Console → Profile → Auth Tokens). For production: use Instance Principal or Resource Principal (no hardcoded credentials, tokens auto-rotate). IAM policies control who can produce/consume from which streams. Enable/disable access through OCI IAM policies without touching application code.

❓ "How do you handle OCI Streaming's 7-day retention limit for use cases needing longer history?"

Use Service Connector Hub to automatically archive all messages from OCI Streaming to OCI Object Storage in real time. Object Storage has no retention limit (keep forever) and costs far less than active streaming storage. Consumers read from the live 7-day stream window. For replaying data older than 7 days, read from Object Storage. For very long-term event sourcing, consider OCI Autonomous Data Warehouse or GoldenGate as the long-term store.

❓ "What is the difference between OCI Streaming REST API and Kafka API?"

Kafka API (port 9092): Uses Kafka binary wire protocol. Works with any Kafka client library. Supports consumer groups, offset management, all standard Kafka operations. Auth via SASL. REST API (HTTPS): Uses OCI HTTP API. No Kafka client library needed. Auth via OCI API key signature. Reading uses cursor-based model (not integer offsets). No native consumer groups (manage cursor manually). Values are base64-encoded. Best for: OCI Functions, quick integrations, languages without Kafka clients.

❓ "How do you migrate an existing Kafka application to OCI Streaming?"

Three steps: (1) Create Stream Pool and Streams in OCI matching your existing topic names. (2) Generate OCI Auth Token for your user (or configure Instance Principal for production). (3) Update application config: change bootstrap_servers to OCI endpoint, add SASL_SSL security protocol, set SASL PLAIN mechanism, set username as tenancy/user, password as auth token. Application business logic code: unchanged. Kafka client library: unchanged. Consumer group IDs: unchanged. Test in dev, validate message flow, switch production. Typical migration: <1 day engineering effort.


🎉 Final Summary 

☁️ OCI Streaming = Managed Kafka: Oracle-managed service that implements the Kafka wire protocol. Zero broker management. Automatic HA across Availability Domains.
🔄 1:1 Concept Mapping: Stream Pool = Kafka Cluster, Stream = Topic, Partition = Partition, Stream Group = Consumer Group, Offset = Offset. Identical semantics.
🔌 Zero Code Migration: Change only bootstrap_servers and SASL authentication config. All producer/consumer/Kafka Streams/Kafka Connect code unchanged.
🔐 IAM Authentication: SASL_SSL with tenancy/username + Auth Token. Production: Instance Principal (no hardcoded secrets, auto-rotating credentials).
⏱️ 7-Day Max Retention: Set all streams to 168 hours. Use Service Connector Hub → Object Storage for unlimited long-term archival at low cost.
📊 Throughput Model: 1 MB/s write, 2 MB/s read per partition. 25 partitions max = 25 MB/s total write. Scale by adding partitions.
🌐 Two API Surfaces: Kafka protocol (port 9092, full Kafka API) and OCI REST API (HTTPS, cursor-based, no Kafka client needed). Use Kafka API for existing apps. REST for serverless/functions.
❌ Key Limitations vs Full Kafka: No Kafka Transactions (exactly-once), no log compaction, no Kafka Admin API (use OCI Console/SDK), 7-day max retention, 25 partition default limit.
🔗 Native OCI Integrations: Service Connector Hub → Object Storage, Functions, Notifications, Logging. OCI GoldenGate for CDC. All zero-code via console.
💡 Choose OCI Streaming when: You're on OCI, want Kafka without ops overhead, need native OCI integrations, moderate scale (≤25 MB/s total), and ≤7-day retention is sufficient.
✅ The Most Important Insight:

OCI Streaming is not trying to be "like Kafka" — it IS Kafka, delivered as an Oracle Cloud service with full Kafka API compatibility and zero infrastructure management.

The engineering decision is straightforward: if you are on OCI and need a Kafka-compatible streaming service without the operational burden of running your own cluster — OCI Streaming is the answer. Your existing Kafka applications, consumer groups, offset management, partition strategies, and Kafka client libraries all work, unchanged. 🎯


Happy Building on Oracle Cloud! 🔥

Comments