Apache Kafka
Self-managed, Confluent Cloud, Amazon MSK, Azure Event Hubs (Kafka protocol)
Access mode: Read-only
Trino connector. Queryable via Flume’s Lakehouse. This system can be accessed both through its native protocol (for metadata introspection) and via Trino federation (for data profiling and cross-system analytical queries).
Required information
| Field | Details |
|---|---|
| Bootstrap servers | Comma-separated list of broker addresses (host1:9092,host2:9092). |
| Port | Default 9092 (plaintext), 9093 (SSL), 9094 (SASL_SSL). Confluent Cloud: 9092 with SASL_SSL. |
| Topic(s) | Kafka topics to access. Topic naming patterns if using wildcards. |
| Schema Registry URL | If using Avro/Protobuf/JSON Schema: Schema Registry endpoint (Confluent, Apicurio, or Karapace). |
| Message format | Avro, Protobuf, JSON, JSON Schema, or raw bytes. Determines how Trino maps topics to table schemas. |
| Authentication | SASL/PLAIN, SASL/SCRAM, SSL certificates, or no auth (internal). Confluent Cloud: API key + secret. |
Network considerations
Self-managed: VPN or direct access. All broker addresses must be reachable. The bootstrap list is just for initial discovery; the client connects to every broker in the cluster.
Confluent Cloud: SASL_SSL over port 9092. No VPN needed.
Amazon MSK: VPC-based. Requires VPC peering or PrivateLink. MSK Serverless has VPC-based access only.
Azure Event Hubs: Kafka protocol on port 9093 with SASL_SSL. FQDN: <namespace>.servicebus.windows.net.
Schema Registry: Separate endpoint (default port 8081). Must also be network-reachable.
Trino access: Flume’s Lakehouse queries Kafka topics as tables via the Trino Kafka connector. Schema Registry integration maps topic messages to columns.
Credential and auth management
Confluent Cloud: API key + secret (per-cluster or service account scoped). Granular ACLs via RBAC.
SASL/SCRAM: Username + password stored in ZooKeeper/KRaft. Create a dedicated user with READ on target topics and consumer groups.
SSL certificates: Client cert + key for mTLS. Provide CA cert for trust.
MSK IAM auth: Flume authenticates via IAM role. Requires kafka-cluster:Connect, kafka-cluster:ReadData, kafka-cluster:DescribeTopic permissions.
ACLs: Kafka ACLs are topic-level. Flume needs READ and DESCRIBE on target topics, plus READ on at least one consumer group.
Validation checks
| Check | Method | Expected result |
|---|---|---|
| Network reachability | kafka-broker-api-versions --bootstrap-server <host>:9092 | Returns broker API versions |
| Authentication | kafka-topics --bootstrap-server <host>:9092 --list (with auth config) | Lists topics |
| Topic access | kafka-console-consumer --topic <topic> --max-messages 1 | Consumes one message |
| Schema Registry | curl https://<registry>:8081/subjects | Lists registered schemas |
| Trino query | SELECT * FROM kafka.<topic> LIMIT 1 via Trino | Returns a row |
Every connection starts from the pre-engagement checklist and goes through the universal validation protocol before production sign-off.