You need to enable JavaScript to run this app.
文档中心
数据库传输服务

数据库传输服务

复制全文
下载 pdf
数据订阅最佳实践
通过 Kafka 消费火山引擎 Proto 格式的订阅数据
复制全文
下载 pdf
通过 Kafka 消费火山引擎 Proto 格式的订阅数据
数据库传输服务 DTS 的数据订阅服务支持使用 Kafka 客户端消费火山引擎 Proto 格式的订阅数据。本文以订阅云数据库 MySQL 版实例为例,介绍如何使用 Go、Java 和 Python 语言消费 Canal 格式的数据。
前提条件
  • 已安装 protoc,建议使用 protoc 3.18 或以上版本。
说明
您可以执行 protoc -version 查看 protoc 版本。
  • 用于订阅消费数据的客户端需要指定服务端 Kafka 版本号,版本号需为 2.2.x(例如 2.2.2)。您可以在示例代码中指定 Kafka 版本号,具体参数如下表所示。
  • 运行语言
    说明
    Go
    通过代码示例中参数 config.Version 指定服务端 Kafka 版本号。
    Python
    通过示例代码中参数 api_version 指定服务端 Kafka 版本号。
    Java
    通过 maven pom.xml 文件中参数 version 指定服务端 Kafka 版本号。
  • 按需安装运行语言环境。
  • 运行语言
    说明
    Go
    安装 Go,需使用 Go 1.13 或以上版本。您可以执行 go version 查看 Go 的版本。
    Python
    1. 安装 Python,需使用 Python 2.7 或以上版本。您可以执行 python --version 查看 Python 的版本。
    1. 依次执行以下命令,安装 pip 依赖。
    • pip install kafka-python
    • pip install protobuf
    • pip install python-snappy
    Java
    1. 安装 Java,需使用 Java 1.8 或以上版本。您可以执行 java -version 查看 Java 版本。
    1. 安装 maven,需使用 Maven 3.8 或以上版本。 您可以执行 mvn -version 查看 Maven 版本。
    1. 在 IDEA 软件,单击 Create New Project 创建一个 Project。
    1. 在新建的 Project 中的项目对象模型文件 pom.xml 中添加以下依赖,本示例以 Kafka 2.2.2 版本为例。同时,您也可以将 pom.xml 文件中 kafka-clients 的版本修改为其他版本 。
    • <dependency>
      <groupId>org.apache.kafka</groupId>
      <artifactId>kafka-clients</artifactId>
      <version>2.2.2</version>
      </dependency>
      <dependency>
      <groupId>com.github.daniel-shuy</groupId>
      <artifactId>kafka-protobuf-serde</artifactId>
      <version>2.2.0</version>
      </dependency>
      <dependency>
      <groupId>org.xerial.snappy</groupId>
      <artifactId>snappy-java</artifactId>
      <version>1.1.8.4</version>
      </dependency>
      <dependency>
      <groupId>com.google.protobuf</groupId>
      <artifactId>protobuf-java</artifactId>
      <version>3.22.2</version>
      </dependency>
关联 Kafka 和订阅任务
本文以 macOS 操作系统为例,介绍如何关联 Kafka 和订阅任务。
  1. 登录 DTS 控制台,创建并配置数据订阅通道。详细信息,请参见订阅方案概览
  1. 在目标数据订阅通道中新增消费组。详细信息,请参见新建消费组
  1. 按需选择 Java 消费示例Python 消费示例,Python 语言和 Java 语言各消费示例的目录如下所示:
  • const list = [
    {
    "lang": "Python 语言",
    "text": `.
    ├── dts_kafka_consumer_demo.py # 消费 Demo 文件
    ├── volc.proto # 火山引擎格式文件
    └── volc_pb2.py # 编译 Volc.proto 后的生成的 Python 文件
    `},
    {
    "lang": "Java 语言",
    "text": `.
    ├── DTSKafkaConsumerDemo.java # 消费 Demo 文件
    ├── Volc.java # 编译 Volc.proto 后的生成的 Java 文件
    └── Volc.proto # 火山引擎格式文件
    `,
    },
    ]
    return (<PreCodeTabs list={list} />);
说明
Go 语言中仅包含一个 Demo 文件。
  1. 按需修改 Demo 文件,具体代码如下所示。
  • 参数
    说明
    示例值
    GROUP
    消费组名称。
    285fef6b91754d0bbaab32e4976c****:test_dtssdk
    USER
    Kafka 用户名。
    test_user
    PASSWORD
    Kafka 用户密码。
    Test@Pwd
    TOPIC
    目标 DTS 数据订阅通道的 Topic。
    d73e98e7fa9340faa3a0d4ccfa10****
    BROKERS
    目标 DTS 数据订阅通道的私网地址。
    kafka-cndvhw9ves******.kafka.ivolces.com:9092
  • const list = [
    {
    "lang": "Go 语言",
    "text": `package main
    import (
    "context"
    "fmt"
    "log"
    "strings"
    "sync"
    "github.com/Shopify/sarama"
    proto "github.com/volcengine/volc-sdk-golang/example/dts/data-subscription-demo/proto"
    protobuf "google.golang.org/protobuf/proto"
    )
    type Handler struct {
    topic string
    partitionCount map[int32]int
    totalCount int
    mu sync.Mutex
    }
    type Config struct {
    username string
    password string
    topic string
    group string
    brokers string
    }
    var (
    c Config
    )
    func init() {
    c.brokers = "your brokers addrress"
    c.topic = "your topic"
    c.group = "your group"
    c.username = "your username"
    c.password = "your password"
    }
    func (h *Handler) Setup(session sarama.ConsumerGroupSession) error {
    fmt.Println("setup")
    return nil
    }
    func (h *Handler) Cleanup(sarama.ConsumerGroupSession) error {
    fmt.Println("clean up")
    return nil
    }
    func (h *Handler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
    fmt.Println("ConsumeClaim")
    for m := range claim.Messages() {
    h.handleMsg(m)
    session.MarkMessage(m, "")
    session.Commit()
    }
    return nil
    }
    func (h *Handler) handleMsg(msg *sarama.ConsumerMessage) {
    h.mu.Lock()
    defer h.mu.Unlock()
    h.totalCount++
    h.partitionCount[msg.Partition]++
    entry := &proto.Entry{}
    if err := protobuf.Unmarshal(msg.Value, entry); err != nil {
    panic(err)
    }
    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)
    }
    func main() {
    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
    topic := c.topic
    group := c.group
    addr := strings.Split(c.brokers, ",")
    cons, err := sarama.NewConsumerGroup(addr, group, config)
    if err != nil {
    panic(err)
    }
    defer cons.Close()
    handler := &Handler{
    topic: topic,
    partitionCount: make(map[int32]int),
    }
    for {
    err = cons.Consume(context.Background(), []string{handler.topic}, handler)
    if err != nil {
    log.Fatalln(err)
    }
    }
    }
    `,
    "selected": true,
    },
    {
    "lang": "Python 语言",
    "text": `
    from kafka import KafkaConsumer
    import volc_pb2
    if __name__ == '__main__':
    brokers = "your brokers address"
    topic = "your topic"
    group = "your group"
    username = "your username"
    password = "your password"
    # create consumer
    consumer = KafkaConsumer(
    topic,
    group_id=group,
    # init consume offset
    auto_offset_reset='latest',
    enable_auto_commit=True,
    # set up SASL authentication
    security_protocol="SASL_PLAINTEXT",
    sasl_mechanism="PLAIN",
    sasl_plain_username=username,
    sasl_plain_password=password,
    api_version=(2, 2, 2),
    bootstrap_servers=brokers.split(','))
    for message in consumer:
    msg = volc_pb2.Entry()
    msg.ParseFromString(message.value)
    print(msg)
    `,
    "selected": true,
    },
    {
    "lang": "Java 语言",
    "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) {
    this.topic = topic;
    // 配置 sasl 认证
    String jaasTemplate = "org.apache.kafka.common.security.scram.ScramLoginModule required username=\"%s\" password=\"%s\";";
    String jaasCfg = String.format(jaasTemplate, username, password);
    // 配置 kafka 参数
    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 {
    // init kafka consumer
    KafkaConsumer<String, Volc.Entry> consumer =
    new KafkaConsumer<String, Volc.Entry>(props,new StringDeserializer(),new KafkaProtobufDeserializer<>(Volc.Entry.parser()));
    // subscribe topic
    consumer.subscribe(Arrays.asList(topic));
    // Loop consuming messages
    while (true) {
    ConsumerRecords<String, Volc.Entry> records = consumer.poll(Duration.ofSeconds(1));
    for (ConsumerRecord<String, Volc.Entry> record : records) {
    // simply printed here
    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);
    // start consume
    c.consume();
    }
    }
    `,
    },
    ]
    return (<PreCodeTabs list={list} />);
运行测试
说明
  • 本文以数据库 test,表格 demo 为例。
  • 根据目标语言选择合适的 JSON 数据。
  1. 在源数据库中,执行以下命令创建一张名为 demo 的表。
  • CREATE TABLE demo (id_t INT);
  • 预期输出:
  • const list = [
    {
    "lang": "Go语言和Python语言",
    "text": `
    src_type:MySQL
    entry_type:DDL
    timestamp:1639057424
    server_id:"105198****"
    database:"test"
    table:"demo"
    ddl_event:{sql:"create table demo (id_t int)"}`,
    "selected": true,
    },
    {
    "lang": "Java语言",
    "text": `
    ConsumerRecord(topic = d73e98e7fa9340faa3a0d4ccfa10****, partition = 0, leaderEpoch = 0, offset = 117, CreateTime = 1639120364271, serialized key size = 9, serialized value size = 67, headers = RecordHeaders(headers = [], isReadOnly = false), key = test.demo, value = src_type: MySQL
    entry_type: DDL
    timestamp: 1639120364
    server_id: "1051985825"
    database: "test"
    table: "demo"
    ddl_event {
    sql: "create table demo (Id_t int)"
    }
    )
    `,
    },
    ]
    return (<PreCodeTabs list={list} />);
  1. 执行以下命令,在 demo 表中插入一条数据:
  • INSERT INTO demo (id_t) VALUES (1);
  • 预期输出:
  • const list = [
    {
    "lang": "Go 语言和 Python 语言",
    "text": `
    src_type:MySQL
    entry_type:DML
    timestamp:1639057434
    server_id:"105198****"
    database:"test"
    table:"demo"
    dml_event:{
    type:INSERT
    table_id:"148"
    column_defs:{
    index:1
    type:INTEGER
    OriginType:"int"
    name:"id_t"
    is_nullable:true
    is_unsigned:false
    }
    rows:{
    after_cols:{
    is_null:false
    int64_value:1
    }
    }
    }`,
    "selected": true,
    },
    {
    "lang": "Java语言",
    "text": `
    ConsumerRecord(topic = d73e98e7fa9340faa3a0d4ccfa10****, partition = 0, leaderEpoch = 0, offset = 118, CreateTime = 1639120462650, serialized key size = 9, serialized value size = 73, headers = RecordHeaders(headers = [], isReadOnly = false), key = test.demo, value = src_type: MySQL
    entry_type: DML
    timestamp: 1639120462
    server_id: "1051985825"
    database: "test"
    table: "demo"
    dml_event {
    type: INSERT
    table_id: "149"
    column_defs {
    index: 1
    type: INTEGER
    OriginType: "int"
    name: "Id_t"
    is_nullable: true
    is_unsigned: false
    }
    rows {
    after_cols {
    is_null: false
    int64_value: 1
    }
    }
    }
    )`,
    },
    ]
    return (<PreCodeTabs list={list} />);
最近更新时间:2026.09.02 11:04:03
这个页面对您有帮助吗?
有用
有用
无用
无用