Kafka Stream API
Overview
The Kafka Stream API delivers measurement and device-change events from Haltian IoT to your Kafka consumer. Device measurements, device-group measurements, and location positions are all delivered as measurement events. Each API key linked to an active Kafka integration receives its own topic. The production broker is available to customers through its public SASL endpoint.
How It Works
When an API key has a live token and is linked to an active Kafka integration, Haltian IoT provisions a topic for the key and forwards supported events to it.
The customer connection uses the API key’s JWT with SASL/OAUTHBEARER over TLS. A separate topic is created for each linked API key.
Set Up Access
Your user needs the Manager role to manage Kafka integrations and API keys. Authenticate to the production Service API and include the returned access token as a Bearer token in each GraphQL request. See Service API Authentication.
Production GraphQL endpoint: https://haltian-iot-api.eu.haltian.io/v1/graphql.
1. Create or use an active Kafka integration
Each organization can have one Kafka integration. If your organization does not already have one, contact Haltian Support to enable Kafka before creating it. The integration must be active.
mutation CreateKafkaIntegration {
createKafkaIntegration(object: {
description: "Production event stream"
kafkaIntegrationState: "active"
}) {
id
kafkaIntegrationState
}
}
2. Create an API key
Create a key for the consumer application. The API returns the JWT only when the key is created, so store it securely. Set expirationDays from 1 to 365 and use the returned expiresAt value to plan token refresh.
mutation CreateApiKey {
createApiKey(createApiKeyInput: {
name: "Production Kafka consumer"
expirationDays: 365
}) {
id
name
organizationId
token {
token
issuedAt
expiresAt
}
}
}
Use the returned organizationId as your organization ID when configuring the consumer group ID.
3. Link the key to the integration
Use assignKafkaIntegrationApiKey to link the API key to the active integration. A valid token and this link enable Kafka provisioning; a separate Kafka role assignment is not required.
mutation LinkApiKeyToKafkaIntegration($integrationId: uuid!, $apiKeyId: uuid!) {
assignKafkaIntegrationApiKey(object: {
kafkaIntegrationId: $integrationId
apiKeyId: $apiKeyId
}) {
kafkaIntegrationId
apiKeyId
}
}
Connect to Kafka
Configure your Kafka client with these settings:
| Setting | Value |
|---|---|
| Bootstrap server | public-eu-haltian-prod.h.aivencloud.com:16036 |
| Security protocol | SASL_SSL |
| SASL mechanism | OAUTHBEARER |
| CA certificate | Request the CA certificate from Haltian and save it as ca.pem |
| OAuth token | The API key JWT returned when you created the key |
| Consumer group ID | Must start with the organizationId returned by createApiKey, followed by a hyphen: <org-id>- |
Provide the API key JWT as the OAuth bearer token. Client libraries expose token callbacks differently; the Python example below uses confluent-kafka.
Find Your Topic
The production topic is derived from the API key ID: haltian-iot.events.apikey.<key-id>. Use the id returned by createApiKey. The topic name is not currently returned through GraphQL.
Events and Payloads
Device measurements, device-group measurements, and location positions use kind: "measurement". Device changes use kind: "entity_change".
kind | Payload fields |
|---|---|
measurement | Device measurements, device-group measurements, and location positions. Fields vary by source and can include type, value, measured_at, and originating_device |
entity_change | Device changes with entity_type, entity_id, change_type, and changed_at |
The Kafka record key is the organization ID. The organization_id field is omitted from the payload because the topic is scoped to one organization. For entity-change events, the extra field is also omitted. Measurement values and optional fields depend on the source event.
Measurement event
{
"kind": "measurement",
"type": "<measurement-type>",
"value": "<measurement-value>",
"measured_at": "<ISO-8601 timestamp>",
"originating_device": {
"id": "<device-uuid>",
"device_model": "<device-model>"
}
}
Device change event
{
"kind": "entity_change",
"entity_type": "device",
"entity_id": "<device-uuid>",
"change_type": "<change-type>",
"changed_at": "<ISO-8601 timestamp>"
}
Retention and Consumer Groups
Topics are currently configured with 7-day retention. Use a stable consumer group ID to resume from its committed offsets. The group ID must begin with <org-id>-; use different IDs when separate applications need independent progress.
Python Consumer Example
Install the client:
pip install confluent-kafka
Save the CA certificate requested from Haltian as ca.pem. Set HALTIAN_KAFKA_TOKEN to the API key JWT, HALTIAN_KAFKA_TOKEN_EXPIRES_AT to its expiresAt value, HALTIAN_API_KEY_ID to the API key ID, and HALTIAN_ORG_ID to the organizationId returned by createApiKey.
import json
import os
from datetime import datetime
from confluent_kafka import Consumer
token = os.environ["HALTIAN_KAFKA_TOKEN"]
expires_at = datetime.fromisoformat(
os.environ["HALTIAN_KAFKA_TOKEN_EXPIRES_AT"].replace("Z", "+00:00")
).timestamp()
api_key_id = os.environ["HALTIAN_API_KEY_ID"]
organization_id = os.environ["HALTIAN_ORG_ID"]
def oauth_cb(_oauth_config):
return token, expires_at
consumer = Consumer({
"bootstrap.servers": "public-eu-haltian-prod.h.aivencloud.com:16036",
"security.protocol": "SASL_SSL",
"ssl.ca.location": "ca.pem",
"sasl.mechanisms": "OAUTHBEARER",
"oauth_cb": oauth_cb,
"group.id": f"{organization_id}-my-consumer",
"auto.offset.reset": "earliest",
})
topic = f"haltian-iot.events.apikey.{api_key_id}"
consumer.subscribe([topic])
try:
while True:
message = consumer.poll(1.0)
if message is None:
continue
if message.error():
print(f"Consumer error: {message.error()}")
continue
print(json.dumps(json.loads(message.value()), indent=2))
finally:
consumer.close()
auto.offset.reset: earliest starts at the earliest retained event when the consumer group has no committed offset. Existing groups resume from their committed offsets. Refresh the API key before it expires. Refreshing immediately invalidates the old JWT, so update and restart the consumer with the new JWT straight away.
Manage Access
- Refresh the API key token through the Service API key mutations. The topic remains the same; the broker access control is updated for the new token.
- Unlinking the key from its Kafka integration, deleting the key, or deleting the integration removes the key’s topic and access controls when it is no longer eligible.
- Setting the integration inactive or in error stops forwarding while preserving the topic and access controls. Reactivating it resumes the stream.
Troubleshooting
| Symptom | Check |
|---|---|
| Authentication or authorization failure | Use the API key JWT with SASL_SSL and OAUTHBEARER. Confirm the key has a live token and is linked to an active integration. |
| Topic not found | Confirm the API key ID used to derive the topic, and that the integration is active and the key is linked. |
| No events arrive | Confirm the organization has live source events and that the linked integration is active. |