xk6_kafka_rest

package module
v0.1.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Apr 26, 2026 License: Apache-2.0 Imports: 13 Imported by: 0

README

xk6-kafka-rest

A k6 extension for publishing JSON messages to Confluent Kafka via the Confluent REST Proxy with full OAuth 2.0 (Client Credentials) support.

⚠️ This extension is currently under active development. APIs may change between versions. See Current Limitations for features not yet supported.


Table of Contents


Features

  • ✅ OAuth 2.0 Client Credentials grant — automatic token fetch & refresh
  • ✅ JSON batch produce (POST /topics/{topic})
  • ✅ Configurable batch size limit (default 500, hard ceiling 1000)
  • ✅ Per-record error detection — surfaces silent Kafka-level failures
  • ✅ Custom k6 metrics: kafka_rest_messages_sent, kafka_rest_publish_duration, kafka_rest_publish_errors
  • ✅ Per-topic metric tags
  • ✅ Custom message keys

Requirements

  • Go 1.21+
  • xk6 — go install go.k6.io/xk6/cmd/xk6@latest

Current Limitations

This extension is under active development. The following features are not yet supported:

Feature Notes
Message headers The REST Proxy v2 JSON API does not support per-record headers
Avro / Schema Registry JSON payloads only at this time

Build

xk6 build --with github.com/Waleed2660/xk6-kafka-rest@latest

# This produces ./k6 — use it instead of the system k6 binary
./k6 version

Quick Start

import { KafkaRestClient } from 'k6/x/kafka-rest';

const client = new KafkaRestClient({
  baseUrl:      __ENV.KAFKA_REST_URL,       // e.g. https://xxx.confluent.cloud
  tokenUrl:     __ENV.OAUTH_TOKEN_URL,      // e.g. https://idp.example.com/oauth/token
  clientId:     __ENV.CLIENT_ID,
  clientSecret: __ENV.CLIENT_SECRET,
  scope:        'kafka',                    // optional
});

export default function () {
  const result = client.produce('my-topic', [
    { key: 'order-123', value: { event: 'ORDER_PLACED', amount: 99.99 } },
  ]);
  console.log(JSON.stringify(result));
}

Run it:

KAFKA_REST_URL=https://... \
OAUTH_TOKEN_URL=https://... \
CLIENT_ID=xxx \
CLIENT_SECRET=yyy \
./k6 run script.js

Local Development (Docker)

Spin up a full local stack — Kafka, REST Proxy, mock OAuth server, and Kafbat UI:

cd local-dev
docker compose up -d
# UI available at http://localhost:8090

Then run the bundled test script:

./k6 run local-dev/test-script.js

Examples

See the /examples folder for ready-to-run scripts:

Script What it shows
single-message.js One message per iteration with check() assertions
batch-produce.js 50 K messages in 200-record batches across 10 VUs
multi-topic.js Multiple topics with independent per-topic thresholds

API Reference

new KafkaRestClient(config)
Option Type Default Description
baseUrl string required Confluent REST Proxy base URL
tokenUrl string required OAuth 2.0 token endpoint
clientId string required OAuth client ID
clientSecret string required OAuth client secret
scope string "" OAuth scope (space-separated)
maxBatchSize number 500 Max records per produce() call. Hard ceiling: 1000

client.produce(topic, messages)

Publishes a batch of messages to a Kafka topic in a single HTTP request.

Parameters

Name Type Description
topic string Kafka topic name
messages Message[] Array of message objects

Message shape

{
  key?:  string | object   // optional partition key
  value: object            // message payload (serialised as JSON)
}

Returns ProduceResponse

{
  offsets: {
    partition:   number
    offset:      number
    error_code?: number   // non-zero means this record was rejected by Kafka
    error?:      string
  }[]
}

⚠️ Per-record errors: The REST Proxy can return HTTP 200 OK while individual records have failed (non-zero error_code). The extension detects these and returns an error with details, also counting them in kafka_rest_publish_errors.


client.close()

No-op. Reserved for future connection cleanup.


Metrics

Metric Type Tags Description
kafka_rest_messages_sent Counter topic Number of records confirmed by Kafka
kafka_rest_publish_duration Trend (ms) topic Round-trip time of each produce HTTP call
kafka_rest_publish_errors Counter topic Number of failed records (HTTP or per-record)

Use them in thresholds:

export const options = {
  thresholds: {
    'kafka_rest_publish_duration{topic:my-topic}': ['p(95)<500'],
    'kafka_rest_publish_errors':                   ['count==0'],
  },
};

IDE Autocomplete (TypeScript)

Type definitions are included in index.d.ts. To enable autocomplete in VS Code or any TypeScript-aware editor:

1. Copy or symlink the types into your project

# from your k6 scripts folder
cp /path/to/xk6-kafka-rest/index.d.ts ./xk6-kafka-rest.d.ts

2. Add a path mapping to tsconfig.json

{
  "compilerOptions": {
    "paths": {
      "k6/x/kafka-rest": ["./xk6-kafka-rest.d.ts"]
    }
  }
}

You'll now get full autocomplete and type checking for KafkaRestClient, Message, ProduceResponse, and all config options.


License

Apache 2.0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type ClientConfig

type ClientConfig struct {
	BaseURL      string `js:"baseUrl"`
	ClientID     string `js:"clientId"`
	ClientSecret string `js:"clientSecret"`
	TokenURL     string `js:"tokenUrl"`
	Scope        string `js:"scope"`
	MaxBatchSize int    `js:"maxBatchSize"`
}

ClientConfig holds all options passed from the JS constructor.

type KafkaRestClient

type KafkaRestClient struct {
	// contains filtered or unexported fields
}

KafkaRestClient is the object exposed to k6 JS scripts.

func (*KafkaRestClient) Close

func (c *KafkaRestClient) Close()

Close is a no-op for now but satisfies the expected JS API surface.

func (*KafkaRestClient) Produce

func (c *KafkaRestClient) Produce(topic string, messages []Message) (*ProduceResponse, error)

Produce publishes a batch of messages to the given Kafka topic. JS usage: client.produce("my-topic", [{value: {foo: "bar"}}])

type KafkaRestModule

type KafkaRestModule struct {
	// contains filtered or unexported fields
}

func (*KafkaRestModule) Exports

func (m *KafkaRestModule) Exports() modules.Exports

type Message

type Message struct {
	Key   interface{} `json:"key,omitempty"`
	Value interface{} `json:"value"`
}

Message represents a single Kafka record sent to the REST Proxy.

type OffsetMetadata

type OffsetMetadata struct {
	Partition int    `json:"partition"`
	Offset    int64  `json:"offset"`
	ErrorCode *int   `json:"error_code,omitempty"`
	Error     string `json:"error,omitempty"`
}

OffsetMetadata describes where a single record landed in Kafka.

type ProduceResponse

type ProduceResponse struct {
	KeySchemaID   int              `json:"key_schema_id,omitempty"`
	ValueSchemaID int              `json:"value_schema_id,omitempty"`
	Offsets       []OffsetMetadata `json:"offsets"`
}

ProduceResponse is returned to the JS script after a successful publish.

type RootModule

type RootModule struct{}

func (*RootModule) NewModuleInstance

func (*RootModule) NewModuleInstance(vu modules.VU) modules.Instance

type TokenManager

type TokenManager struct {
	// contains filtered or unexported fields
}

TokenManager fetches and caches an OAuth client-credentials token. It is safe for concurrent use across goroutines (VUs).

func NewTokenManager

func NewTokenManager(cfg ClientConfig) *TokenManager

NewTokenManager creates a TokenManager from a ClientConfig.

func (*TokenManager) Token

func (t *TokenManager) Token(ctx context.Context) (string, error)

Token returns a valid Bearer token, refreshing it when fewer than 30 seconds remain before expiry.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL