Which Message Queue Systems Does DBX Support for Admin Functionality?

DBX currently provides full admin functionality for Apache Pulsar, backend-only support for Apache Kafka (requiring a Java agent), and an unimplemented placeholder for Apache RocketMQ that returns a "not yet implemented" error.

DBX is an open-source database management platform that includes a pluggable message queue (MQ) admin console for managing streaming infrastructure. Understanding which message queue systems DBX supports for admin functionality helps platform teams evaluate whether it fits their Kafka, Pulsar, or RocketMQ environments. The system implements a trait-based adapter pattern that separates generic admin operations from vendor-specific implementations.

Supported Message Queue Systems

DBX implements a pluggable architecture centered on the MessageQueueAdmin trait and the MqAdminRegistry cache. The registry lazily builds adapters in build_adapter based on the MqAdminConfig.system_kind field defined in crates/dbx-core/src/mq/types.rs.

Apache Pulsar (Fully Supported)

Apache Pulsar is the only system with complete administrative coverage in DBX. The PulsarAdmin adapter in crates/dbx-core/src/mq/adapters/pulsar.rs implements all tenant, namespace, topic, subscription, policy, and raw API operations.

When MqAdminRegistry::build_adapter receives a configuration with system_kind: MqSystemKind::Pulsar, it instantiates the full adapter that exposes capabilities including tenant management and namespace creation.

Apache Kafka (Backend Ready, UI Pending)

The Kafka adapter exists at crates/dbx-core/src/mq/adapters/kafka.rs and supports many admin operations including topics, consumer groups, ACLs, and retention policies. However, the UI does not expose Kafka options yet.

The implementation requires an AgentLaunchSpec to spawn a Java agent for communicating with Kafka brokers. According to the source in crates/dbx-core/src/mq/mod.rs lines 45-50, the registry reserves a slot for Kafka and will use KafkaAdmin when an agent launch spec is provided.

Apache RocketMQ (Reserved but Not Implemented)

RocketMQ is reserved in the type system but lacks implementation. The MqAdminRegistry::build_adapter method in crates/dbx-core/src/mq/mod.rs at line 151 returns Err("RocketMQ admin is not yet implemented") when attempting to create a RocketMQ adapter.

Working with DBX Message Queue Adapters

Testing a Pulsar Connection

The frontend can test Pulsar connectivity using the mqTestConnection function, which calls the backend service defined in crates/dbx-core/src/mq/service.rs:

import { mqTestConnection } from '@/lib/mq-api'

const cfg = {
  systemKind: 'pulsar',
  adminUrl: 'https://pulsar.example.com:8443',
  auth: { kind: 'token', token: 'MY_JWT_TOKEN' },
  tlsSkipVerify: false,
}

// Returns a `MqClusterInfo` with capabilities (tenants, namespaces, etc.)
const info = await mqTestConnection(cfg)
console.log(info.capabilities)   // true for all Pulsar features

Configuring the Kafka Admin Adapter

For Kafka, you must provide an AgentLaunchSpec to instantiate the KafkaAdmin struct. The adapter entry point is KafkaAdmin::new at lines 52-66 of crates/dbx-core/src/mq/adapters/kafka.rs:

use dbx_core::mq::MqAdminConfig;
use dbx_core::mq::adapters::kafka::KafkaAdmin;
use dbx_core::mq::port::MessageQueueAdmin;

// Build a config that points at a Kafka broker cluster
let cfg = MqAdminConfig {
    system_kind: MqSystemKind::Kafka,
    admin_url: String::new(),               // not used for Kafka
    auth: MqAuth::None,
    tls_skip_verify: false,
    extra: serde_json::json!({ "bootstrapServers": "localhost:9092" }),
    ..Default::default()
};

// Launch spec for the Java agent (must be supplied by the host)
let launch = AgentLaunchSpec::new("java", vec!["-jar", "kafka-agent.jar"]);
let admin = KafkaAdmin::new(cfg, launch).await?;

// List topics (Kafka only supports topics, not tenants/namespaces)
let topics = admin.list_topics(&NamespaceRef::default(), ListTopicsOpts::default()).await?;
println!("Kafka topics: {:?}", topics);

Handling RocketMQ Placeholder Errors

Attempting to configure RocketMQ results in the expected error from the registry:

let cfg = MqAdminConfig {
    system_kind: MqSystemKind::RocketMq,
    ..Default::default()
};
let registry = MqAdminRegistry::new();
match registry.build_transient_config(cfg, None).await {
    Ok(_) => println!("Unexpected – RocketMQ not implemented"),
    Err(e) => println!("RocketMQ admin error: {}", e), // → "RocketMQ admin is not yet implemented"
}

This error originates from the match arm in build_adapter at line 151 of crates/dbx-core/src/mq/mod.rs.

Summary

Frequently Asked Questions

Does DBX support Kafka message queue administration?

DBX includes a functional KafkaAdmin adapter that supports topics, consumer groups, ACLs, and retention policies, but the UI does not yet expose Kafka configuration options. You can use the Rust API directly by providing an AgentLaunchSpec to spawn the required Java agent as shown in crates/dbx-core/src/mq/adapters/kafka.rs.

What Pulsar operations are available in DBX?

The PulsarAdmin adapter provides comprehensive coverage including tenant management, namespace operations, topic administration, subscription handling, policy enforcement, and raw API access. All capabilities are exposed through the frontend UI and validated via the mqTestConnection RPC handler in crates/dbx-core/src/mq/service.rs.

Why does RocketMQ return an error in DBX?

RocketMQ is reserved in the MqSystemKind enum defined in crates/dbx-core/src/mq/types.rs, but the implementation in MqAdminRegistry::build_adapter at line 151 of crates/dbx-core/src/mq/mod.rs explicitly returns an error stating "RocketMQ admin is not yet implemented."

How does DBX handle message queue authentication?

DBX supports token-based authentication for Pulsar through the auth field in MqAdminConfig. For Kafka, authentication parameters are passed via the extra JSON field, though specific mechanisms depend on the Java agent configuration specified in AgentLaunchSpec.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →