"github.com/Shopify/sarama"
proto "github.com/volcengine/volc-sdk-golang/example/dts/data-subscription-demo/proto"
protobuf "google.golang.org/protobuf/proto"
partitionCount map[int32]int
c.brokers = "your brokers addrress"
c.username = "your username"
c.password = "your password"
func (h *Handler) Setup(session sarama.ConsumerGroupSession) error {
func (h *Handler) Cleanup(sarama.ConsumerGroupSession) error {
func (h *Handler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
fmt.Println("ConsumeClaim")
for m := range claim.Messages() {
session.MarkMessage(m, "")
func (h *Handler) handleMsg(msg *sarama.ConsumerMessage) {
h.partitionCount[msg.Partition]++
if err := protobuf.Unmarshal(msg.Value, entry); err != nil {
fmt.Println("-------------- handle message --------------")
fmt.Printf("get message EventType:%v\n", entry.EntryType.String())
switch entry.GetEntryType() {
case proto.EntryType_DDL:
event := entry.GetDdlEvent()
fmt.Printf("ddl %v\n", event.Sql)
case proto.EntryType_DML:
event := entry.GetDmlEvent()
cols := event.ColumnDefs
for _, row := range event.Rows {
var before, after []string
for i, col := range row.BeforeCols {
before = append(before, fmt.Sprintf("%+v[%+v]", cols[i].GetName(), col.GetValue()))
for i, col := range row.AfterCols {
after = append(after, fmt.Sprintf("%+v[%+v]", cols[i].GetName(), col.GetValue()))
fmt.Printf("get row before=%v after=%v\n", before, after)
fmt.Printf("fetch message partition=%v key=%v\n", msg.Partition, string(msg.Key))
fmt.Printf("count partition-count=%v total-count=%v\n", h.partitionCount, h.totalCount)
fmt.Printf("config: %+v", c)
config := sarama.NewConfig()
config.Net.SASL.User = c.username
config.Net.SASL.Password = c.password
config.Net.SASL.Enable = true
config.Consumer.Offsets.Initial = sarama.OffsetNewest
config.Version = sarama.V2_2_2_0
addr := strings.Split(c.brokers, ",")
cons, err := sarama.NewConsumerGroup(addr, group, config)
partitionCount: make(map[int32]int),
err = cons.Consume(context.Background(), []string{handler.topic}, handler)
from kafka import KafkaConsumer
if __name__ == '__main__':
brokers = "your brokers address"
username = "your username"
password = "your password"
consumer = KafkaConsumer(
auto_offset_reset='latest',
enable_auto_commit=True,
# set up SASL authentication
security_protocol="SASL_PLAINTEXT",
sasl_plain_username=username,
sasl_plain_password=password,
bootstrap_servers=brokers.split(','))
for message in consumer:
msg.ParseFromString(message.value)
"text": `package dts.sub.volc;
import com.github.daniel.shuy.kafka.protobuf.serde.KafkaProtobufDeserializer;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.Arrays;
import java.util.Properties;
public class DTSKafkaConsumerDemo {
private final String topic;
private final Properties props;
public DTSKafkaConsumerDemo(String brokers, String topic, String group, String username, String password) {
String jaasTemplate = "org.apache.kafka.common.security.scram.ScramLoginModule required username=\"%s\" password=\"%s\";";
String jaasCfg = String.format(jaasTemplate, username, password);
props = new Properties();
props.put("bootstrap.servers", brokers);
props.put("group.id", group);
props.put("enable.auto.commit", "true");
props.put("auto.commit.interval.ms", "1000");
props.put("auto.offset.reset", "latest");
props.put("session.timeout.ms", "30000");
props.put("security.protocol", "SASL_PLAINTEXT");
props.put("sasl.mechanism", "PLAIN");
props.put("sasl.jaas.config", jaasCfg);
public void consume() throws java.lang.Exception {
KafkaConsumer<String, Volc.Entry> consumer =
new KafkaConsumer<String, Volc.Entry>(props,new StringDeserializer(),new KafkaProtobufDeserializer<>(Volc.Entry.parser()));
consumer.subscribe(Arrays.asList(topic));
// Loop consuming messages
ConsumerRecords<String, Volc.Entry> records = consumer.poll(Duration.ofSeconds(1));
for (ConsumerRecord<String, Volc.Entry> record : records) {
System.out.println(record);
public static void main(String[] args)throws java.lang.Exception {
// set your kafka brokers address
String brokers = "your brokers address";
// set your kafka username and password
String username = "yout username";
String password = "your password";
// set your kafka topic 和 group
String topic = "your topic";
String group = "your group";
DTSKafkaConsumerDemo c = new DTSKafkaConsumerDemo(brokers,topic,group, username, password);
return (<PreCodeTabs list={list} />);