Kafka Stream API

Connect to Haltian IoT Kafka in production and consume real-time measurements and device events

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.

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:

SettingValue
Bootstrap serverpublic-eu-haltian-prod.h.aivencloud.com:16036
Security protocolSASL_SSL
SASL mechanismOAUTHBEARER
CA certificateRequest the CA certificate from Haltian and save it as ca.pem
OAuth tokenThe API key JWT returned when you created the key
Consumer group IDMust 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".

kindPayload fields
measurementDevice measurements, device-group measurements, and location positions. Fields vary by source and can include type, value, measured_at, and originating_device
entity_changeDevice 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

SymptomCheck
Authentication or authorization failureUse 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 foundConfirm the API key ID used to derive the topic, and that the integration is active and the key is linked.
No events arriveConfirm the organization has live source events and that the linked integration is active.