Subscribe to Channel via WebSocket

Subscribe to a Data Streams channel to receive data from publishers in real time over a WebSocket connection.

Before you begin, ensure you have a channel created, access control configured, and a token pair available. Subscriptions are typically established from client-side applications (web or mobile) because WebSocket connections are long-lived and must remain open to receive data continuously.

Note: The JavaScript SDK supports browser environments for subscriber operations.

Install the Data Streams SDK

copy
npm install @zcatalyst/datastreams

Add to pom.xml:

copy


    com.zoho.catalyst
    datastreams
    LATEST

copy
pip install zcatalyst-sdk

Get Token Pair from Server

First, request a token pair from your server endpoint (created in generate-token-pair):

copy

// Client-side example
const response = await fetch('https://your-server.com/get-token', {
    method: 'POST',
    headers: { 'Content-Type': 'application/json' },
    body: JSON.stringify({
        channelId: 'YOUR_CHANNEL_ID',
        userId: 'current_user_id'
    })
});

const { token } = await response.json(); // token contains: { url, wss_id, key, channel_id }

Create WebSocket Connection

Use the token pair to create a WebSocket client instance. The SDK handles authentication automatically.

copy
const { DataStreamsWebSocket } = require('@zcatalyst/datastreams');

const socket = new DataStreamsWebSocket({
    url: token.url,
    zuid: token.wss_id,
    key: token.key,
    enableLogging: true  // Optional: enable debug logging
});
copy
import com.zc.component.datastream.ZCDatastream;
import com.zc.component.datastream.beans.ZCDataStreamsConfig;
import com.zc.component.datastream.websocket.ZCDataStreamsWebSocket;

ZCDatastream datastream = ZCDatastream.getInstance();

ZCDataStreamsConfig config = new ZCDataStreamsConfig(
    tokenResponse.getUrl(),
    tokenResponse.getWssId(),
    tokenResponse.getKey(),
    false  // enableLogging
);

ZCDataStreamsWebSocket socket = datastream.createWebSocketClient(config);
copy
from zcatalyst_sdk.datastreams.websocket_client import DataStreamsWebSocket
from zcatalyst_sdk.types.datastreams import DataStreamsConfig

config = DataStreamsConfig(
    url=token["url"],
    zuid=token["wss-id"],
    key=token["key"],
    enable_logging=True
)

socket = DataStreamsWebSocket(config=config)
Note: In Python, the WebSocket connection is not opened until you call socket.connect() after registering your event listeners.

Choose Your Subscriber Type

Before subscribing, decide which subscriber type fits your use case:

Type Value When to Use
Live events only ‘0’ Start receiving only NEW messages published after subscribing
Earliest available ‘-1’ Receive ALL available messages from the earliest stored event (within 48-hour retention)
Resume ‘-2’ Resume from where you last left off (or live events for new subscribers)
Specific Streaming ID ‘’ Resume from a specific message (use the streamingId from a previous message)
Note: Data retention is 48 hours. Events older than 48 hours are not available.

Register Event Listeners and Subscribe

copy
// Listen for connection opened
socket.on('open', () => {
    console.log('Connected! Subscribing to channel...');
    
    // Subscribe with your chosen subscriber type
    socket.subscribe('0');  // Live events only
    // socket.subscribe('-1');  // Earliest available
    // socket.subscribe('-2');  // Resume
    // socket.subscribe('16965000000027481');  // From specific streaming ID
});

// Listen for incoming messages
socket.on('message', (event) => {
    console.log('Received:', event.data);
    console.log('Streaming ID:', event.streamingId);
    
    // Process your message here
    // ...
    
    // CRITICAL: Always send acknowledgement to receive next message
    socket.sendAck();
});

// Listen for errors
socket.on('error', (err) => {
    console.error('Error:', err.code, err.message);
});

// Listen for connection closed
socket.on('close', () => {
    console.log('Connection closed');
});
copy
import com.zc.component.datastream.websocket.ZCDataStreamsEventListener;
import com.zc.component.datastream.beans.ZCCustomEvent;
import com.zc.component.datastream.beans.ZCDataStreamMessageEvent;

socket.addEventListener(new ZCDataStreamsEventListener() {
    @Override
    public void onOpen(ZCCustomEvent event) {
        System.out.println("Connected! Subscribing to channel...");
        
        // Subscribe with your chosen subscriber type
        socket.subscribe("0");  // Live events only
        // socket.subscribe("-1");  // Earliest available
        // socket.subscribe("-2");  // Resume
    }
    
    @Override
    public void onMessage(ZCDataStreamMessageEvent event) {
        System.out.println("Received: " + event.getData());
        System.out.println("Streaming ID: " + event.getStreamingId());
        
        // Process your message here
        // ...
        
        // CRITICAL: Always send acknowledgement to receive next message
        socket.sendAck();
    }
    
    @Override
    public void onError(ZCCustomEvent event) {
        System.err.println("Error: " + event.getMessage());
    }
    
    @Override
    public void onClose(ZCCustomEvent event) {
        System.out.println("Connection closed");
    }
    
    @Override
    public void onPong(ZCCustomEvent event) {
        // Keep-alive pong received
    }
});
copy
def on_open(event):
    print("Connected! Subscribing to channel...")
    
    # Subscribe with your chosen subscriber type
    socket.subscribe("0")  # Live events only
    # socket.subscribe("-1")  # Earliest available
    # socket.subscribe("-2")  # Resume

def on_message(event):
    print("Received:", event.data)
    print("Streaming ID:", event.streaming_id)
    
    # Process your message here
    # ...
    
    # CRITICAL: Always send acknowledgement to receive next message
    socket.send_ack()

def on_error(event):
    print("Error:", event.code, event.message)

def on_close(event):
    print("Connection closed")

# Register event listeners
socket.set_on_open(on_open)
socket.set_on_message(on_message)
socket.set_on_error(on_error)
socket.set_on_close(on_close)

# Connect (Python requires explicit connect call)
socket.connect()

Handle Message Events

Each message event received contains:

Field Type Description
operation string Event type: ’event' for inline data or ‘api’ for bulk data
streamingId string Unique identifier for this message in the stream
data string The published data (for inline events)
url string URL to fetch bulk data (for bulk/api events)
method string HTTP method for fetching bulk data (for bulk/api events)

Example: Processing Different Event Types

copy

socket.on('message', (event) => {
    if (event.operation === 'event') {
        // Inline data
        console.log('Data:', event.data);
        const parsedData = JSON.parse(event.data);
        // Process parsedData...
    } else if (event.operation === 'api') {
        // Bulk data - fetch from URL
        console.log('Bulk data available at:', event.url);
        // Fetch and process bulk data...
    }
// Always acknowledge
socket.sendAck();

});

Send Acknowledgement

You must call sendAck() after processing each message.

  • Until acknowledgement is sent, the next message will NOT be delivered
  • This ensures ordered, reliable message delivery
  • If you disconnect before acknowledging, you can reconnect using subscriber type ‘-2’ to resume
copy

// JavaScript
socket.sendAck();

// Java socket.sendAck();

// Python socket.send_ack()

Unsubscribe and Disconnect

When you’re done receiving data:

copy
// Unsubscribe: ends the session. Generate a new token pair to subscribe again.
socket.unsubscribe();

// Or disconnect only: reconnect later and resume with subscribe type "-2".
socket.close();
copy
// Unsubscribe: ends the session. Generate a new token pair to subscribe again.
socket.unsubscribe();

// Or disconnect only: reconnect later and resume with subscribe type "-2".
socket.close();
copy
# Unsubscribe: ends the session. Generate a new token pair to subscribe again.
socket.unsubscribe()

# Or disconnect only: reconnect later and resume with subscribe type "-2".
socket.close()

Example Implementation

copy

const { DataStreamsWebSocket } = require('@zcatalyst/datastreams');

async function subscribeToChannel() { // 1. Get token from server const response = await fetch(‘https://your-server.com/get-token', { method: ‘POST’, body: JSON.stringify({ channelId: ‘YOUR_CHANNEL_ID’, userId: ‘user123’ }) }); const { token } = await response.json();

// 2. Create WebSocket connection
const socket = new DataStreamsWebSocket({
    url: token.url,
    zuid: token.wss_id,
    key: token.key,
    enableLogging: true
});

// 3. Register event listeners
socket.on('open', () => {
    console.log('Connected!');
    socket.subscribe('0');  // Live events
});

socket.on('message', (event) => {
    console.log('Received:', event.data);
    
    // Process message
    const data = JSON.parse(event.data);
    updateUI(data);
    
    // Acknowledge
    socket.sendAck();
});

socket.on('error', (err) => {
    console.error('Error:', err);
});

socket.on('close', () => {
    console.log('Disconnected');
});

}

subscribeToChannel();

The SDK manages the connection lifecycle automatically: opens the WebSocket connection, validates the token pair, subscribes when you call subscribe(), and pushes messages with the acknowledgement flow.

You can find further details on Subscriber Types, Message Event Structure, Acknowledgement, Connection Lifecycle, and Data Retention.

Last Updated 2026-10-05 20:43:57 +0530 IST