curl --request POST \
--url https://api.streamkap.com/topics/bulk/snapshot \
--header 'Authorization: Bearer <token>' \
--header 'Content-Type: application/json' \
--data '
{
"topic_ids": [
"<string>"
],
"snapshot_type": "incremental"
}
'import requests
url = "https://api.streamkap.com/topics/bulk/snapshot"
payload = {
"topic_ids": ["<string>"],
"snapshot_type": "incremental"
}
headers = {
"Authorization": "Bearer <token>",
"Content-Type": "application/json"
}
response = requests.post(url, json=payload, headers=headers)
print(response.text)const options = {
method: 'POST',
headers: {Authorization: 'Bearer <token>', 'Content-Type': 'application/json'},
body: JSON.stringify({topic_ids: ['<string>'], snapshot_type: 'incremental'})
};
fetch('https://api.streamkap.com/topics/bulk/snapshot', options)
.then(res => res.json())
.then(res => console.log(res))
.catch(err => console.error(err));<?php
$curl = curl_init();
curl_setopt_array($curl, [
CURLOPT_URL => "https://api.streamkap.com/topics/bulk/snapshot",
CURLOPT_RETURNTRANSFER => true,
CURLOPT_ENCODING => "",
CURLOPT_MAXREDIRS => 10,
CURLOPT_TIMEOUT => 30,
CURLOPT_HTTP_VERSION => CURL_HTTP_VERSION_1_1,
CURLOPT_CUSTOMREQUEST => "POST",
CURLOPT_POSTFIELDS => json_encode([
'topic_ids' => [
'<string>'
],
'snapshot_type' => 'incremental'
]),
CURLOPT_HTTPHEADER => [
"Authorization: Bearer <token>",
"Content-Type: application/json"
],
]);
$response = curl_exec($curl);
$err = curl_error($curl);
curl_close($curl);
if ($err) {
echo "cURL Error #:" . $err;
} else {
echo $response;
}package main
import (
"fmt"
"strings"
"net/http"
"io"
)
func main() {
url := "https://api.streamkap.com/topics/bulk/snapshot"
payload := strings.NewReader("{\n \"topic_ids\": [\n \"<string>\"\n ],\n \"snapshot_type\": \"incremental\"\n}")
req, _ := http.NewRequest("POST", url, payload)
req.Header.Add("Authorization", "Bearer <token>")
req.Header.Add("Content-Type", "application/json")
res, _ := http.DefaultClient.Do(req)
defer res.Body.Close()
body, _ := io.ReadAll(res.Body)
fmt.Println(string(body))
}HttpResponse<String> response = Unirest.post("https://api.streamkap.com/topics/bulk/snapshot")
.header("Authorization", "Bearer <token>")
.header("Content-Type", "application/json")
.body("{\n \"topic_ids\": [\n \"<string>\"\n ],\n \"snapshot_type\": \"incremental\"\n}")
.asString();require 'uri'
require 'net/http'
url = URI("https://api.streamkap.com/topics/bulk/snapshot")
http = Net::HTTP.new(url.host, url.port)
http.use_ssl = true
request = Net::HTTP::Post.new(url)
request["Authorization"] = 'Bearer <token>'
request["Content-Type"] = 'application/json'
request.body = "{\n \"topic_ids\": [\n \"<string>\"\n ],\n \"snapshot_type\": \"incremental\"\n}"
response = http.request(request)
puts response.read_body{
"results": [
{
"id": "<string>",
"success": true,
"message": "<string>",
"data": {}
}
],
"summary": {
"total": 123,
"succeeded": 123,
"failed": 123
}
}{
"detail": [
{
"loc": [
"<string>"
],
"msg": "<string>",
"type": "<string>",
"input": "<unknown>",
"ctx": {}
}
]
}Bulk Snapshot Topics
Execute snapshot for selected topics, grouped by their source connector.
Accepts topic DB IDs, resolves them to their source connectors, and triggers per-source snapshots scoped to only the selected topics (via topic_names).
Only topics produced by source connectors can be snapshotted - topics from transforms, destinations, or manually created topics are skipped.
Supports both incremental (default) and blocking snapshot types:
- incremental: Uses watermarking to capture data while streaming continues.
- blocking: Pauses streaming during snapshot. Required for keyless tables.
Results are reported per source connector, since snapshots are source-level operations. Returns partial success results - continues processing even if individual snapshots fail.
curl --request POST \
--url https://api.streamkap.com/topics/bulk/snapshot \
--header 'Authorization: Bearer <token>' \
--header 'Content-Type: application/json' \
--data '
{
"topic_ids": [
"<string>"
],
"snapshot_type": "incremental"
}
'import requests
url = "https://api.streamkap.com/topics/bulk/snapshot"
payload = {
"topic_ids": ["<string>"],
"snapshot_type": "incremental"
}
headers = {
"Authorization": "Bearer <token>",
"Content-Type": "application/json"
}
response = requests.post(url, json=payload, headers=headers)
print(response.text)const options = {
method: 'POST',
headers: {Authorization: 'Bearer <token>', 'Content-Type': 'application/json'},
body: JSON.stringify({topic_ids: ['<string>'], snapshot_type: 'incremental'})
};
fetch('https://api.streamkap.com/topics/bulk/snapshot', options)
.then(res => res.json())
.then(res => console.log(res))
.catch(err => console.error(err));<?php
$curl = curl_init();
curl_setopt_array($curl, [
CURLOPT_URL => "https://api.streamkap.com/topics/bulk/snapshot",
CURLOPT_RETURNTRANSFER => true,
CURLOPT_ENCODING => "",
CURLOPT_MAXREDIRS => 10,
CURLOPT_TIMEOUT => 30,
CURLOPT_HTTP_VERSION => CURL_HTTP_VERSION_1_1,
CURLOPT_CUSTOMREQUEST => "POST",
CURLOPT_POSTFIELDS => json_encode([
'topic_ids' => [
'<string>'
],
'snapshot_type' => 'incremental'
]),
CURLOPT_HTTPHEADER => [
"Authorization: Bearer <token>",
"Content-Type: application/json"
],
]);
$response = curl_exec($curl);
$err = curl_error($curl);
curl_close($curl);
if ($err) {
echo "cURL Error #:" . $err;
} else {
echo $response;
}package main
import (
"fmt"
"strings"
"net/http"
"io"
)
func main() {
url := "https://api.streamkap.com/topics/bulk/snapshot"
payload := strings.NewReader("{\n \"topic_ids\": [\n \"<string>\"\n ],\n \"snapshot_type\": \"incremental\"\n}")
req, _ := http.NewRequest("POST", url, payload)
req.Header.Add("Authorization", "Bearer <token>")
req.Header.Add("Content-Type", "application/json")
res, _ := http.DefaultClient.Do(req)
defer res.Body.Close()
body, _ := io.ReadAll(res.Body)
fmt.Println(string(body))
}HttpResponse<String> response = Unirest.post("https://api.streamkap.com/topics/bulk/snapshot")
.header("Authorization", "Bearer <token>")
.header("Content-Type", "application/json")
.body("{\n \"topic_ids\": [\n \"<string>\"\n ],\n \"snapshot_type\": \"incremental\"\n}")
.asString();require 'uri'
require 'net/http'
url = URI("https://api.streamkap.com/topics/bulk/snapshot")
http = Net::HTTP.new(url.host, url.port)
http.use_ssl = true
request = Net::HTTP::Post.new(url)
request["Authorization"] = 'Bearer <token>'
request["Content-Type"] = 'application/json'
request.body = "{\n \"topic_ids\": [\n \"<string>\"\n ],\n \"snapshot_type\": \"incremental\"\n}"
response = http.request(request)
puts response.read_body{
"results": [
{
"id": "<string>",
"success": true,
"message": "<string>",
"data": {}
}
],
"summary": {
"total": 123,
"succeeded": 123,
"failed": 123
}
}{
"detail": [
{
"loc": [
"<string>"
],
"msg": "<string>",
"type": "<string>",
"input": "<unknown>",
"ctx": {}
}
]
}Authorizations
Bearer authentication header of the form Bearer <token>, where <token> is your auth token.
Body
Request model for bulk snapshot operations at the topic level.
Accepts topic DB IDs, groups them by their source connector, and triggers per-source snapshots with the selected topics only.
List of topic DB IDs (_id) to operate on.
1 - 500 elementsType of snapshot to execute. 'incremental' (default) uses watermarking while streaming continues. 'blocking' pauses streaming during snapshot - required for keyless tables.
incremental, blocking Was this page helpful?