Documentation
¶
Overview ¶
+marmot:name=Kafka +marmot:description=This plugin discovers Kafka topics from Kafka clusters. +marmot:status=experimental
Index ¶
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type AuthConfig ¶
type AuthConfig struct {
Type string `json:"type" yaml:"type" description:"Authentication type: none, sasl_plaintext, sasl_ssl, ssl"`
Username string `json:"username,omitempty" yaml:"username,omitempty" description:"SASL username"`
Password string `json:"password,omitempty" yaml:"password,omitempty" description:"SASL password"`
Mechanism string `json:"mechanism,omitempty" yaml:"mechanism,omitempty" description:"SASL mechanism: PLAIN, SCRAM-SHA-256, SCRAM-SHA-512"`
TLSCertPath string `json:"tls_cert_path,omitempty" yaml:"tls_cert_path,omitempty" description:"Path to TLS certificate file"`
TLSKeyPath string `json:"tls_key_path,omitempty" yaml:"tls_key_path,omitempty" description:"Path to TLS key file"`
TLSCACertPath string `json:"tls_ca_cert_path,omitempty" yaml:"tls_ca_cert_path,omitempty" description:"Path to TLS CA certificate file"`
TLSSkipVerify bool `json:"tls_skip_verify,omitempty" yaml:"tls_skip_verify,omitempty" description:"Skip TLS verification"`
}
Authentication configuration
type Config ¶
type Config struct {
plugin.BaseConfig `json:",inline"`
// Connection configuration
BootstrapServers string `json:"bootstrap_servers" yaml:"bootstrap_servers" description:"Comma-separated list of bootstrap servers"`
ClientID string `json:"client_id" yaml:"client_id" description:"Client ID for the consumer"`
Authentication *AuthConfig `json:"authentication,omitempty" yaml:"authentication,omitempty" description:"Authentication configuration"`
ConsumerConfig map[string]string `json:"consumer_config,omitempty" yaml:"consumer_config,omitempty" description:"Additional consumer configuration"`
ClientTimeout int `json:"client_timeout_seconds" yaml:"client_timeout_seconds" description:"Request timeout in seconds"`
// Schema Registry configuration (optional)
SchemaRegistry *SchemaRegistryConfig `json:"schema_registry,omitempty" yaml:"schema_registry,omitempty" description:"Schema Registry configuration"`
// Topic patterns for filtering
TopicFilter *plugin.Filter `json:"topic_filter,omitempty" yaml:"topic_filter,omitempty" description:"Filter configuration for topics"`
// Metadata extraction
IncludePartitionInfo bool `` /* 126-byte string literal not displayed */
IncludeTopicConfig bool `json:"include_topic_config" yaml:"include_topic_config" description:"Whether to include topic configuration in metadata"`
}
Config for Kafka plugin +marmot:config
type ConsumerGroupDetails ¶
type ConsumerGroupDetails struct {
State string
Protocol string
ProtocolType string
Members []ConsumerGroupMember
}
ConsumerGroupDetails holds information about a Kafka consumer group
type ConsumerGroupMember ¶
type ConsumerGroupMember struct {
ClientID string
ClientHost string
TopicPartitions map[string][]int
}
ConsumerGroupMember holds information about a member of a consumer group
type KafkaConsumerGroupFields ¶
type KafkaConsumerGroupFields struct {
GroupId string `json:"group_id" metadata:"group_id" description:"Consumer group ID"`
State string `json:"state" metadata:"state" description:"Current state of the consumer group"`
Protocol string `json:"protocol" metadata:"protocol" description:"Rebalance protocol"`
ProtocolType string `json:"protocol_type" metadata:"protocol_type" description:"Protocol type"`
SubscribedTopics []string `json:"subscribed_topics" metadata:"subscribed_topics" description:"Topics the group is subscribed to"`
Members []string `json:"members" metadata:"members" description:"Members of the consumer group"`
}
KafkaConsumerGroupFields represents Kafka consumer group-specific metadata fields +marmot:metadata
type KafkaTopicFields ¶
type KafkaTopicFields struct {
TopicName string `json:"topic_name" metadata:"topic_name" description:"Name of the Kafka topic"`
PartitionCount int32 `json:"partition_count" metadata:"partition_count" description:"Number of partitions"`
ReplicationFactor int16 `json:"replication_factor" metadata:"replication_factor" description:"Replication factor"`
RetentionMs string `json:"retention_ms" metadata:"retention.ms" description:"Message retention period in milliseconds"`
RetentionBytes string `json:"retention_bytes" metadata:"retention.bytes" description:"Maximum size of the topic in bytes"`
CleanupPolicy string `json:"cleanup_policy" metadata:"cleanup.policy" description:"Topic cleanup policy"`
MinInsyncReplicas string `json:"min_insync_replicas" metadata:"min.insync.replicas" description:"Minimum number of in-sync replicas"`
MaxMessageBytes string `json:"max_message_bytes" metadata:"max.message.bytes" description:"Maximum message size in bytes"`
SegmentBytes string `json:"segment_bytes" metadata:"segment.bytes" description:"Segment file size in bytes"`
SegmentMs string `json:"segment_ms" metadata:"segment.ms" description:"Segment file roll time in milliseconds"`
DeleteRetentionMs string `json:"delete_retention_ms" metadata:"delete.retention.ms" description:"Time to retain deleted segments in milliseconds"`
ValueSchemaId int `json:"value_schema_id" metadata:"value_schema_id" description:"ID of the value schema in Schema Registry"`
ValueSchemaVersion int `json:"value_schema_version" metadata:"value_schema_version" description:"Version of the value schema"`
ValueSchemaType string `json:"value_schema_type" metadata:"value_schema_type" description:"Type of the value schema (AVRO, JSON, etc.)"`
ValueSchema string `json:"value_schema" metadata:"value_schema" description:"Value schema definition"`
KeySchemaId int `json:"key_schema_id" metadata:"key_schema_id" description:"ID of the key schema in Schema Registry"`
KeySchemaVersion int `json:"key_schema_version" metadata:"key_schema_version" description:"Version of the key schema"`
KeySchemaType string `json:"key_schema_type" metadata:"key_schema_type" description:"Type of the key schema (AVRO, JSON, etc.)"`
KeySchema string `json:"key_schema" metadata:"key_schema" description:"Key schema definition"`
}
KafkaTopicFields represents Kafka topic-specific metadata fields +marmot:metadata
type SchemaRegistryConfig ¶
type SchemaRegistryConfig struct {
URL string `json:"url" yaml:"url" description:"Schema Registry URL"`
Config map[string]string `json:"config,omitempty" yaml:"config,omitempty" description:"Additional Schema Registry configuration"`
Enabled bool `json:"enabled" yaml:"enabled" description:"Whether to use Schema Registry"`
}
Schema Registry configuration
type Source ¶
type Source struct {
// contains filtered or unexported fields
}
func (*Source) Discover ¶
func (s *Source) Discover(ctx context.Context, pluginConfig plugin.RawPluginConfig) (*plugin.DiscoveryResult, error)
type TopicDetails ¶
type TopicDetails struct {
// contains filtered or unexported fields
}
TopicDetails holds information about a Kafka topic
Click to show internal directories.
Click to hide internal directories.