首页
学习
活动
专区
圈层
工具
发布
社区首页 >问答首页 >多个Debezium源连接器不能同时工作

多个Debezium源连接器不能同时工作
EN

Stack Overflow用户
提问于 2021-02-10 21:34:37
回答 1查看 445关注 0票数 0

Docker Compose

代码语言:javascript
复制
kafka:
image: confluentinc/cp-enterprise-kafka:6.0.0
container_name: kafka
depends_on:
  - zookeeper
ports:
  - 9092:9092
environment:
  KAFKA_BROKER_ID: 1
  KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
  KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
  KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
  KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://kafka:9092
  KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true"
  KAFKA_METRIC_REPORTERS: io.confluent.metrics.reporter.ConfluentMetricsReporter
  KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
  KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
  KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
  KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 100
  CONFLUENT_METRICS_REPORTER_BOOTSTRAP_SERVERS: kafka:29092
  CONFLUENT_METRICS_REPORTER_ZOOKEEPER_CONNECT: zookeeper:2181
  CONFLUENT_METRICS_REPORTER_TOPIC_REPLICAS: 1
  CONFLUENT_METRICS_ENABLE: 'true'
  CONFLUENT_SUPPORT_CUSTOMER_ID: 'anonymous'
  KAFKA_LOG_RETENTION_MS: 100000000 
  KAFKA_LOG_RETENTION_CHECK_INTERVAL_MS: 5000

连接器1: Debezium源连接器(像这样,我需要8个连接器用于8个表)

代码语言:javascript
复制
curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" localhost:8083/connectors/ -d '
    {
        "name": "mysql5-mdmembers",
        "config": {
            "connector.class": "io.debezium.connector.mysql.MySqlConnector",
            "tasks.max": "10",
            "database.hostname": "13.232.63.40",
            "database.port": "3307",
            "database.user": "root",
            "database.password": "secret",
            "database.server.id": "11",
            "database.server.name": "dbserver",
            "database.whitelist": "indianmo_imc_new",
            "table.whitelist": "indianmo_imc_new.md_members_cdc",
            "database.history.kafka.bootstrap.servers": "kafka:29092",
            "database.history.kafka.topic": "mysql5_md_members",
            "key.converter": "io.confluent.connect.avro.AvroConverter",
            "value.converter": "io.confluent.connect.avro.AvroConverter",
            "key.converter.schema.registry.url": "http://schema-registry:8081",
            "value.converter.schema.registry.url": "http://schema-registry:8081",
            "transforms": "unwrap,dropTopicPrefix,selectFields,renameFields,addTopicPrefix,convertTS",
            "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
            "transforms.dropTopicPrefix.type":"org.apache.kafka.connect.transforms.RegexRouter",
            "transforms.dropTopicPrefix.regex":"dbserver.indianmo_imc_new.(.*)",
            "transforms.dropTopicPrefix.replacement":"$1",
            "transforms.selectFields.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
            "transforms.selectFields.whitelist": "mem_id,mem_mobile,mem_fname,mem_email,mem_gender,mem_dob,mem_marital_status,mem_state,mem_city,mem_zip,mem_primary_lang,mem_created_on,mem_updated_on,mem_tot_ttt,transfer_count,first_transfer_date,last_transfer_date,dnd_status",
            "transforms.renameFields.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
            "transforms.renameFields.renames": "mem_fname:mem_name,mem_primary_lang:mem_primary_language,mem_tot_ttt:mem_total_talktime,transfer_count:mem_transfer_count,first_transfer_date:mem_first_transferred_on,last_transfer_date:mem_last_transferred_on,dnd_status:mem_dnd_status",
            "transforms.addTopicPrefix.type":"org.apache.kafka.connect.transforms.RegexRouter",
            "transforms.addTopicPrefix.regex":"(.*)",
            "transforms.addTopicPrefix.replacement":"mdt_$1",
            "transforms.convertTS.type"       : "org.apache.kafka.connect.transforms.TimestampConverter$Value",
            "transforms.convertTS.field"      : "mem_created_on,first_transfer_date,last_transfer_date",
            "transforms.convertTS.format"     : "YYYY-MM-dd H:mm:ss",
            "transforms.convertTS.target.type": "unix"
        }
    }'

连接器2:同一数据库的Debezium源连接器

代码语言:javascript
复制
curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" localhost:8083/connectors/ -d '
    {
        "name": "mysql5-connector2",
        "config": {
            "connector.class": "io.debezium.connector.mysql.MySqlConnector",
            "tasks.max": "2",
            "database.hostname": "13.232.63.40",
            "database.port": "3307",
            "database.user": "root",
            "database.password": "secret",
            "database.server.id": "11",
            "database.server.name": "dbserver",
            "database.whitelist": "indianmo_imc_new",
            "table.whitelist": "indianmo_imc_new.associate_leads_cdc",
            "database.history.kafka.bootstrap.servers": "kafka:29092",
            "database.history.kafka.topic": "mysql5table",
            "key.converter": "io.confluent.connect.avro.AvroConverter",
            "value.converter": "io.confluent.connect.avro.AvroConverter",
            "key.converter.schema.registry.url": "http://schema-registry:8081",
            "value.converter.schema.registry.url": "http://schema-registry:8081",
            "transforms": "unwrap,dropTopicPrefix,selectFields,renameFields,addTopicPrefix,convertTS",
            "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
            "transforms.dropTopicPrefix.type":"org.apache.kafka.connect.transforms.RegexRouter",
            "transforms.dropTopicPrefix.regex":"dbserver.indianmo_imc_new.(.*)",
            "transforms.dropTopicPrefix.replacement":"$1",
            "transforms.selectFields.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
            "transforms.selectFields.whitelist": "AL_Id,A_Id,L_Id,cityID,rtitle,lead_price,selling_price,lead_trans_type,R_Id,rank,bought_by,bought_ip,leadSentBy,pushed_date,pushed_on,tbSendDate,tbSendOn",
            "transforms.renameFields.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
            "transforms.renameFields.renames": "AL_Id:al_id,A_Id:associate_id,L_Id:mem_id,cityID:city_id,rtitle:product_id,lead_price:lead_actual_price,selling_price:lead_selling_price,lead_trans_type:lead_type_msql5,R_Id:requirement_id,rank:lead_rating,bought_by:lead_bought_by,bought_ip:lead_bought_ip,leadSentBy:lead_sent_by,pushed_date:lead_pushed_date,pushed_on:lead_pushed_on,tbSendDate:lead_sold_date,tbSendOn:lead_sold_on",
            "transforms.addTopicPrefix.type":"org.apache.kafka.connect.transforms.RegexRouter",
            "transforms.addTopicPrefix.regex":"(.*)",
            "transforms.addTopicPrefix.replacement":"ldt_lm_$1",
            "transforms.convertTS.type"       : "org.apache.kafka.connect.transforms.TimestampConverter$Value",
            "transforms.convertTS.field"      : "pushed_on,tbSendOn",
            "transforms.convertTS.format"     : "YYYY-MM-dd H:mm:ss",
            "transforms.convertTS.target.type": "unix"
        }
    }'

接收器连接器(像这样我需要接收8个表)

代码语言:javascript
复制
curl -X PUT http://localhost:8083/connectors/sink-jdbc-mysql-01/config \
    -H "Content-Type: application/json" -d '{
        "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
        "connection.url": "jdbc:mysql://65.0.213.250:3306/demo",
        "topics": "ldt_lm_associate_leads_cdc",
        "key.converter": "io.confluent.connect.avro.AvroConverter",
        "value.converter": "io.confluent.connect.avro.AvroConverter",
        "key.converter.schema.registry.url": "http://schema-registry:8081",
        "value.converter.schema.registry.url": "http://schema-registry:8081",
        "connection.user": "root",
        "connection.password": "secret",
        "auto.create": true,
        "auto.evolve": true,
        "insert.mode": "upsert",
        "delete.enabled": true,
        "pk.mode": "record_key",
        "pk.fields": "al_id",
        "transforms": "RenameKey",
        "transforms.RenameKey.type": "org.apache.kafka.connect.transforms.ReplaceField$Key",
        "transforms.RenameKey.renames": "AL_Id:al_id"
    }'

我从kafka connect得到的错误是

代码语言:javascript
复制
2021-02-10 13:14:45,326] INFO [mysql5-connector2|task-0] Connector task finished all work and is now shutdown (io.debezium.connector.mysql.MySqlConnectorTask:496)

我这里有一个经纪人。如果我使用多个代理,我的问题会得到解决吗?多个代理的docker yml配置是什么?我在源数据库中有多个表。我想把一切都打沉。对于这一点,我想要多个连接器。但每当我运行多个源连接器时,它就会停止前一个连接器(刚刚使用了两个源连接器,它出现了问题,我至少需要8个源连接器和8个宿连接器)。我该怎么做请帮帮忙。提前感谢!

EN

回答 1

Stack Overflow用户

发布于 2021-06-11 18:28:09

票数 0
EN
页面原文内容由Stack Overflow提供。腾讯云小微IT领域专用引擎提供翻译支持
原文链接:

https://stackoverflow.com/questions/66137920

复制
相关文章

相似问题

领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档