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.
Install the Data Streams SDK
npm install @zcatalyst/datastreams
Add to pom.xml:
com.zoho.catalyst
datastreams
LATEST
pip install zcatalyst-sdk
Get Token Pair from Server
First, request a token pair from your server endpoint (created in generate-token-pair):
// 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.
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
});
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);
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)
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) |
Register Event Listeners and Subscribe
// 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');
});
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
}
});
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
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
// JavaScript
socket.sendAck();
// Java
socket.sendAck();
// Python
socket.send_ack()
Unsubscribe and Disconnect
When you’re done receiving data:
// 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();
// 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();
# 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
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
Yes
No
Send your feedback to us