Class KafkaClusterController
java.lang.Object
io.micronaut.controlpanel.panels.kafka.KafkaClusterController
@Controller("${micronaut.control-panel.path:/control-panel}/kafka-control-panel-controller")
@ExecuteOn("blocking")
@Internal
@Requires(beans=org.apache.kafka.clients.admin.AdminClient.class) @Requires(property="micronaut.control-panel.panels.kafka-cluster.enabled",notEquals="false")
public final class KafkaClusterController
extends Object
JSON endpoints used by the Kafka Cluster control panel.
-
Method Summary
Modifier and TypeMethodDescriptionbrokers()consumerGroup(String groupId) ksqldb()messages(String topic, int partition, String mode, @Nullable Long offset, @Nullable Long timestamp, int limit) schemaRegistrySubject(String subject, int version) summary()topicNames(boolean includeInternal) writes()
-
Method Details
-
summary
-
brokers
-
topics
@Get("/topics{?search,includeInternal,start,length}") public KafkaClusterResponse.Section<KafkaClusterResponse.TopicPage> topics(@QueryValue @Nullable String search, @QueryValue(defaultValue="false") boolean includeInternal, @QueryValue(defaultValue="0") int start, @QueryValue(defaultValue="25") int length) -
topicNames
@Get("/topic-names{?includeInternal}") public KafkaClusterResponse.Section<List<String>> topicNames(@QueryValue(defaultValue="false") boolean includeInternal) -
topic
@Get("/topics/{topic}") public KafkaClusterResponse.Section<KafkaClusterResponse.TopicDetail> topic(String topic) -
consumerGroups
@Get("/consumer-groups") public KafkaClusterResponse.Section<List<KafkaClusterResponse.ConsumerGroupSummary>> consumerGroups() -
consumerGroup
@Get("/consumer-groups/{groupId}") public KafkaClusterResponse.Section<KafkaClusterResponse.ConsumerGroupDetail> consumerGroup(String groupId) -
appConsumers
@Get("/app-consumers") public KafkaClusterResponse.Section<List<KafkaClusterResponse.AppConsumer>> appConsumers() -
messages
@Get("/messages{?topic,partition,mode,offset,timestamp,limit}") public KafkaClusterResponse.Section<KafkaClusterResponse.MessagePage> messages(@QueryValue String topic, @QueryValue int partition, @QueryValue(defaultValue="beginning") String mode, @QueryValue @Nullable Long offset, @QueryValue @Nullable Long timestamp, @QueryValue(defaultValue="25") int limit) -
writes
@Get("/writes") public KafkaClusterResponse.Section<KafkaClusterResponse.WriteCapabilities> writes() -
schemaRegistry
@Get("/schema-registry") public KafkaClusterResponse.Section<KafkaClusterResponse.SchemaRegistryOverview> schemaRegistry() -
schemaRegistrySubject
@Get("/schema-registry/subjects/{subject}{?version}") public KafkaClusterResponse.Section<KafkaClusterResponse.SchemaVersionDetail> schemaRegistrySubject(String subject, @QueryValue(defaultValue="1") int version) -
kafkaConnect
@Get("/kafka-connect") public KafkaClusterResponse.Section<KafkaClusterResponse.KafkaConnectOverview> kafkaConnect() -
ksqldb
-
createTopic
@Post("/topics") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> createTopic(@Body KafkaClusterResponse.CreateTopicRequest request) -
updateTopicConfig
@Post("/topics/config") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> updateTopicConfig(@Body KafkaClusterResponse.UpdateTopicConfigRequest request) -
increasePartitions
@Post("/topics/partitions") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> increasePartitions(@Body KafkaClusterResponse.IncreasePartitionsRequest request) -
deleteTopic
@Post("/topics/delete") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> deleteTopic(@Body KafkaClusterResponse.DeleteTopicRequest request) -
produceMessage
@Post("/messages") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> produceMessage(@Body KafkaClusterResponse.ProduceMessageRequest request) -
deleteConsumerGroup
@Post("/consumer-groups/delete") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> deleteConsumerGroup(@Body KafkaClusterResponse.DeleteConsumerGroupRequest request) -
resetConsumerGroupOffsets
@Post("/consumer-groups/reset-offsets") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> resetConsumerGroupOffsets(@Body KafkaClusterResponse.ResetOffsetsRequest request) -
pauseAppConsumer
@Post("/app-consumers/pause") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> pauseAppConsumer(@Body KafkaClusterResponse.AppConsumerActionRequest request) -
resumeAppConsumer
@Post("/app-consumers/resume") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> resumeAppConsumer(@Body KafkaClusterResponse.AppConsumerActionRequest request) -
registerSchema
@Post("/schema-registry/subjects") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> registerSchema(@Body KafkaClusterResponse.RegisterSchemaRequest request) -
updateSchemaCompatibility
@Post("/schema-registry/compatibility") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> updateSchemaCompatibility(@Body KafkaClusterResponse.UpdateSchemaCompatibilityRequest request) -
deleteSchemaSubject
@Post("/schema-registry/subjects/delete") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> deleteSchemaSubject(@Body KafkaClusterResponse.DeleteSchemaSubjectRequest request) -
deleteSchemaVersion
@Post("/schema-registry/subjects/versions/delete") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> deleteSchemaVersion(@Body KafkaClusterResponse.DeleteSchemaVersionRequest request) -
pauseConnector
@Post("/kafka-connect/pause") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> pauseConnector(@Body KafkaClusterResponse.ConnectorActionRequest request) -
resumeConnector
@Post("/kafka-connect/resume") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> resumeConnector(@Body KafkaClusterResponse.ConnectorActionRequest request) -
restartConnector
@Post("/kafka-connect/restart") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> restartConnector(@Body KafkaClusterResponse.ConnectorActionRequest request) -
restartConnectorTask
@Post("/kafka-connect/restart-task") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> restartConnectorTask(@Body KafkaClusterResponse.ConnectorTaskActionRequest request) -
updateConnectorConfig
@Post("/kafka-connect/config") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> updateConnectorConfig(@Body KafkaClusterResponse.UpdateConnectorConfigRequest request) -
deleteConnector
@Post("/kafka-connect/delete") public KafkaClusterResponse.Section<KafkaClusterResponse.ActionResult> deleteConnector(@Body KafkaClusterResponse.ConnectorActionRequest request)
-