# Which Message Queue Systems Does DBX Support for Admin Functionality?

> Discover which message queue systems DBX supports for admin tasks. Learn about Pulsar, Kafka, and RocketMQ integration in t8y2/dbx.

- Repository: [skyler/dbx](https://github.com/t8y2/dbx)
- Tags: how-to-guide
- Published: 2026-07-04

---

**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`](https://github.com/t8y2/dbx/blob/main/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`](https://github.com/t8y2/dbx/blob/main/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`](https://github.com/t8y2/dbx/blob/main/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`](https://github.com/t8y2/dbx/blob/main/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`](https://github.com/t8y2/dbx/blob/main/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`](https://github.com/t8y2/dbx/blob/main/crates/dbx-core/src/mq/service.rs):

```typescript
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`](https://github.com/t8y2/dbx/blob/main/crates/dbx-core/src/mq/adapters/kafka.rs):

```rust
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:

```rust
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`](https://github.com/t8y2/dbx/blob/main/crates/dbx-core/src/mq/mod.rs).

## Summary

- **Apache Pulsar**: Full admin support including tenants, namespaces, topics, and subscriptions via `PulsarAdmin` in [`crates/dbx-core/src/mq/adapters/pulsar.rs`](https://github.com/t8y2/dbx/blob/main/crates/dbx-core/src/mq/adapters/pulsar.rs).
- **Apache Kafka**: Backend adapter complete with topic and consumer group management at [`crates/dbx-core/src/mq/adapters/kafka.rs`](https://github.com/t8y2/dbx/blob/main/crates/dbx-core/src/mq/adapters/kafka.rs), but requires an `AgentLaunchSpec` and remains hidden from the UI pending frontend integration.
- **Apache RocketMQ**: Type system placeholder only; calling `MqAdminRegistry::build_adapter` returns a "not yet implemented" error from [`crates/dbx-core/src/mq/mod.rs`](https://github.com/t8y2/dbx/blob/main/crates/dbx-core/src/mq/mod.rs).

## 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`](https://github.com/t8y2/dbx/blob/main/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`](https://github.com/t8y2/dbx/blob/main/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`](https://github.com/t8y2/dbx/blob/main/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`](https://github.com/t8y2/dbx/blob/main/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`.