import (
"context"
"fmt"
"io"
"cloud.google.com/go/managedkafka/apiv1/managedkafkapb"
"google.golang.org/api/option"
"google.golang.org/protobuf/types/known/fieldmaskpb"
managedkafka "cloud.google.com/go/managedkafka/apiv1"
)
func updateConsumerGroup(w io.Writer, projectID, region, clusterID, consumerGroupID, topicPath string, partitionOffsets map[int32]int64, opts ...option.ClientOption) error {
// projectID := "my-project-id"
// region := "us-central1"
// clusterID := "my-cluster"
// consumerGroupID := "my-consumer-group"
// topicPath := "my-topic-path"
// partitionOffsets := map[int32]int64{1: 10, 2: 20, 3: 30}
ctx := context.Background()
client, err := managedkafka.NewClient(ctx, opts...)
if err != nil {
return fmt.Errorf("managedkafka.NewClient got err: %w", err)
}
defer client.Close()
clusterPath := fmt.Sprintf("projects/%s/locations/%s/clusters/%s", projectID, region, clusterID)
consumerGroupPath := fmt.Sprintf("%s/consumerGroups/%s", clusterPath, consumerGroupID)
partitionMetadata := make(map[int32]*managedkafkapb.ConsumerPartitionMetadata)
for partition, offset := range partitionOffsets {
partitionMetadata[partition] = &managedkafkapb.ConsumerPartitionMetadata{
Offset: offset,
}
}
topicConfig := map[string]*managedkafkapb.ConsumerTopicMetadata{
topicPath: {
Partitions: partitionMetadata,
},
}
consumerGroupConfig := managedkafkapb.ConsumerGroup{
Name: consumerGroupPath,
Topics: topicConfig,
}
paths := []string{"topics"}
updateMask := &fieldmaskpb.FieldMask{
Paths: paths,
}
req := &managedkafkapb.UpdateConsumerGroupRequest{
UpdateMask: updateMask,
ConsumerGroup: &consumerGroupConfig,
}
consumerGroup, err := client.UpdateConsumerGroup(ctx, req)
if err != nil {
return fmt.Errorf("client.UpdateConsumerGroup got err: %w", err)
}
fmt.Fprintf(w, "Updated consumer group: %#v\n", consumerGroup)
return nil
}