标签: confluent-platform

使用kafka接收器重命名elasticsearch中的索引

我正在使用以下水槽。问题是它将elasticsearch索引名称设置为与主题相同。我想要一个不同的 elasticseach 索引名称。我怎样才能做到这一点。我正在使用汇合4

{
  "name": "es-sink-mysql-foobar-02",
  "config": {
    "_comment": "-- standard converter stuff -- this can actually go in the worker config globally --",
    "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "key.converter": "io.confluent.connect.avro.AvroConverter",
    "key.converter.schema.registry.url": "http://localhost:8081",
    "value.converter.schema.registry.url": "http://localhost:8081",


    "_comment": "--- Elasticsearch-specific config ---",
    "_comment": "Elasticsearch server address",
    "connection.url": "http://localhost:9200",

    "_comment": "Elasticsearch mapping name. Gets created automatically if doesn't exist  ",
    "type.name": "type.name=kafka-connect",
    "index.name": "asimtest",
    "_comment": "Which topic to stream data from into Elasticsearch",
    "topics": "mysql-foobar",

    "_comment": "If the Kafka message doesn't have a key …
Run Code Online (Sandbox Code Playgroud)

elasticsearch apache-kafka kafka-consumer-api apache-kafka-connect confluent-platform

1
推荐指数
1
解决办法
3385
查看次数

卡夫卡和 IIDR CDC

我正在尝试使用 DB2--IBM CDC --Kafka构建 CDC 管道 ,并且我正在尝试找出设置此管道的正确方法。我尝试了以下事情 -

1.在Linux on Prem上设置3节点kafka集群

2.使用文件在Linux on prem上安装IIDR CDC软件setup-iidr-11.4.0.1-5085-linux-x86.bin。CDC 实例已启动并正在运行。

各种在线文档建议安装“IIDR 管理控制台”来配置源数据存储和 CDC 服务器配置以及 Kafka 订阅配置来构建管道。

目前我没有安装管理控制台。对此有几个问题 -

1.除了 IBM CDC 管理控制台之外,还有其他替代方案来设置 kafka-CDC 管道吗?

2.如何获取IIDR管理控制台?如果我们将它安装在本地 Windows 桌面上并尝试连接到远程 Linux 服务器上的 CDC/Kafka,它会工作吗?

3.还有其他方法可以将 IIDR CDC 数据摄取到 Kafka 吗?

我对 CDC/ IIDR 还很陌生,请帮忙!

db2 cdc apache-kafka ibm-infosphere confluent-platform

1
推荐指数
1
解决办法
3414
查看次数

无法连接 Kafka 架构注册表中的 localhost:8081

我已经使用以下命令启动了 Zookeeper 和 Kafka:

bin/zookeeper-server-start.sh config/zookeeper.properties
bin/kafka-server-start.sh config/server.properties
Run Code Online (Sandbox Code Playgroud)

当我尝试获取架构注册表兼容性设置(向后、向前、无)时,我运行了以下curl命令:

curl -X GET http://localhost:8081/config
Run Code Online (Sandbox Code Playgroud)

预期的:

{"compatibility":"BACKWARD"}
Run Code Online (Sandbox Code Playgroud)

结果:

curl: (7) Failed to connect to localhost port 8081: Connection refused
Run Code Online (Sandbox Code Playgroud)

我怎样才能找到要使用哪个端口?

apache-kafka confluent-schema-registry confluent-platform

1
推荐指数
1
解决办法
4922
查看次数

如何使用 C# 反序列化 Kafka 中的 Avro 消息

嗨,我正在使用 Confluence kafka。我有返回通用记录的消费者。我想反序列化它。我找不到任何办法。我可以手动完成每个字段,例如

 object options = ((GenericRecord)response.Message.Value["Product"])["Options"];
Run Code Online (Sandbox Code Playgroud)

我在这里找到了一个

使用 C# 反序列化 Avro 文件 但是如何将架构转换为流?我想知道我们是否可以使用任何解决方案反序列化到我们的 c# 模型中?任何帮助将不胜感激。谢谢。

c# avro apache-kafka confluent-schema-registry confluent-platform

1
推荐指数
1
解决办法
9914
查看次数

Kafka-confluent:如何在 JDBC 接收器连接器中使用 pk.mode=record_key 进行更新插入和删除模式?

在 Kafka confluent 中,我们如何在使用pk.mode=record_keyMySQL 表中的复合键的同时使用源作为 CSV 文件的 upsert ?使用pk.mode=record_values. 是否有任何额外的配置需要完成?

如果我尝试使用pk.mode=record_key. 错误 - 由以下原因引起org.apache.kafka.connect.errors.ConnectException::需要定义一个 PK 列,因为记录的关键架构是一种原始类型。以下是我的 JDBC 接收器连接器配置:

    {
    "name": "<name>",
    "config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
    "tasks.max": "1",
    "topics": "<topic name>",
    "connection.url": "<url>",
    "connection.user": "<user name>",
    "connection.password": "*******",
    "insert.mode": "upsert",
    "batch.size": "50000",
    "table.name.format": "<table name>",
    "pk.mode": "record_key",
    "pk.fields": "field1,field2",
    "auto.create": "true",
    "auto.evolve": "true",
    "max.retries": "10",
    "retry.backoff.ms": "3000",
    "mode": "bulk",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter.schemas.enable": "true",
    "value.converter.schema.registry.url": "http://localhost:8081"
  }
}
Run Code Online (Sandbox Code Playgroud)

apache-kafka apache-kafka-connect confluent-platform

1
推荐指数
1
解决办法
2793
查看次数

NotEnoughReplicasException: 当前 ISR Set(2) 的大小不足以满足 min.isr 要求 3

我有以下设置 Brokers : 3 - 所有都启动并运行 min.insync.replicas=3。

我用以下配置创建了一个主题

bin\windows\kafka-topics --zookeeper 127.0.0.1:2181 --topic topic-ack-all --create --partitions 4 --replication-factor 3

我用“ack = all”触发了生产者,生产者能够发送消息。但是,当我启动消费者时问题就开始了

bin\windows\kafka-console-consumer --bootstrap-server localhost:9094,localhost:9092 --topic topic-ack-all --from-beginning

错误是

NotEnoughReplicasException: 当前 ISR Set(2) 的大小不足以满足 3 的 min.isr 要求 NotEnoughReplicasException: 当前 ISR Set(3) 的大小不足以满足分区 __con 的 min.isr 要求 3

我在这里看到两种错误。我浏览了文档,也对“min.isr”有所了解,但是,这些错误消息并不清楚。

  1. 当前 ISR 集是什么意思?每个主题是否不同,它的含义是什么?
  2. 我猜 min.isr 与 min.insync.replicas 相同。我希望应该具有至少与“复制因子”相同的价值?

更新 #1

Topic: topic-ack-all    PartitionCount: 4       ReplicationFactor: 3    Configs:            
        Topic: topic-ack-all    Partition: 0    Leader: 1       Replicas: 1,2,3 Isr: 1,2,3  
        Topic: topic-ack-all    Partition: 1    Leader: …
Run Code Online (Sandbox Code Playgroud)

apache-kafka confluent-platform

1
推荐指数
1
解决办法
3013
查看次数

如何从使用 librdkafka.redist 作为依赖项的 C# 应用程序访问 librdkafka.redist 日志?

我的应用程序控制台日志显示了以下很多内容:

%3|1602097315.970|FAIL|rdkafka#consumer-2| [thrd:kfkqaapq0002d.ch.me.com:9092/bootstrap]: kfkqaapq0002d.ch.me.com:9092/bootstrap: Failed to resolve 'kfkqaapq0002d.ch.me.com:9092': No such host is known.  (after 42ms in state CONNECT)
%3|1602097315.970|FAIL|rdkafka#consumer-3| [thrd:kfkqaapq0002d.ch.me.com:9092/bootstrap]: kfkqaapq0002d.ch.me.com:9092/bootstrap: Failed to resolve 'kfkqaapq0002d.ch.me.com:9092': No such host is known.  (after 41ms in state CONNECT)
%3|1602097315.972|FAIL|rdkafka#producer-1| [thrd:kfkqaapq0003d.ch.me.com:9092/bootstrap]: kfkqaapq0003d.ch.me.com:9092/bootstrap: Failed to resolve 'kfkqaapq0003d.ch.me.com:9092': No such host is known.  (after 48ms in state CONNECT)
%3|1602097315.973|ERROR|rdkafka#producer-1| [thrd:app]: rdkafka#producer-1: kfkqaapq0003d.ch.me.com:9092/bootstrap: Failed to resolve 'kfkqaapq0003d.ch.me.com:9092': No such host is known.  (after 48ms in state CONNECT)
%3|1602097316.459|FAIL|rdkafka#producer-1| [thrd:kfkqaapq0001d.ch.me.com:9092/bootstrap]: kfkqaapq0001d.ch.me.com:9092/bootstrap: Failed to resolve 'kfkqaapq0001d.ch.me.com:9092': No such host …
Run Code Online (Sandbox Code Playgroud)

.net c# apache-kafka .net-core confluent-platform

1
推荐指数
1
解决办法
1226
查看次数

Kafka connect docker image - 无法找到任何实现 Connector 且名称与 ElasticsearchSinkConnector 匹配的类

我已经使用 kafka-connect 图像confluenceinc/cp-kafka-connect一段时间了。根据 Confluence 文档,这个 docker 镜像附带了预安装的连接器插件,包括 Elastic。

我以前一直使用5.4.1-ccs效果很好的版本,我可以添加弹性接收器连接器配置,它们工作得很好。但是我尝试更新 confluentinc/cp-kafka-connect到最新版本v6.0.1,但现在出现错误。

ConnectException: Failed to find any class that implements Connector and which name matches ElasticsearchSinkConnector

我在 Confluence 网站上阅读了很多文档,但有点零星。我知道问题是插件未安装,因为它们被从新的 docker 映像中删除或路径错误(不确定是哪一个)。

我该如何解决这个问题?(注意:我也编写了自己的java插件,所以两者都需要工作)

这是我docker-compose目前的文件(同样,这适用于 version 5.4.1-ccs

kafka-connect-node-1:
  image: confluentinc/cp-kafka-connect:5.4.1 #using old version because of breaking change
  hostname: kafka-connect-node-1
  ports:
    - '8083:8083'
  environment:
    CONNECT_BOOTSTRAP_SERVERS: [MY_SERVER]
    CONNECT_REST_PORT: 8083
    CONNECT_GROUP_ID: compose-connect-group
    CONNECT_CONFIG_STORAGE_TOPIC: connect-configs
    CONNECT_OFFSET_STORAGE_TOPIC: connect-offsets
    CONNECT_STATUS_STORAGE_TOPIC: connect-status
    CONNECT_KEY_CONVERTER: io.confluent.connect.avro.AvroConverter
    CONNECT_KEY_CONVERTER_SCHEMA_REGISTRY_URL: 'http://kafka-schema-registry:8084'
    CONNECT_VALUE_CONVERTER: io.confluent.connect.avro.AvroConverter
    CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_URL: 'http://kafka-schema-registry:8084'
    CONNECT_INTERNAL_KEY_CONVERTER: …
Run Code Online (Sandbox Code Playgroud)

apache-kafka docker apache-kafka-connect confluent-platform

1
推荐指数
1
解决办法
1806
查看次数

基于Kafka的Confluence Platform 7.1是免费的吗?开源?用于生产用途

我有开始使用 Kafka 的用例,并且正在寻找开源免费(生产)kafka。

当检查Confluence 7.1平台看起来很合适时,因为它捆绑了zookeeper/kafka/schema注册表/kafka UI。

在决定继续之前,只想检查一下 Confluence Platform 7.1 是否免费且开源?我需要购买许可或付费支持吗?

apache-kafka confluent-platform

1
推荐指数
1
解决办法
810
查看次数

Kafka Connect JDBC 接收器-topics.regex 不起作用

我正在使用这个debezium-examples

我在jdbc-sink.json 中添加了"topics.regex": "CID1122.(.*)"如下

{
"name": "jdbc-sink",
"config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
    "tasks.max": "1",
    "topics.regex": "CID1122.(.*)",
    "connection.url": "jdbc:mysql://mysql:3306/inventory?verifyServerCertificate=false",
    "connection.user": "root",
    "connection.password": "debezium",
    "auto.create": "true",
    "transforms": "unwrap",
    "transforms.unwrap.type": "io.debezium.transforms.UnwrapFromEnvelope",
    "name": "jdbc-sink",
    "insert.mode": "upsert",
    "pk.fields": "id,companyId",
    "pk.mode": "record_value"
}
}
Run Code Online (Sandbox Code Playgroud)

Kafka 主题列表是

CID1122.department
CID1122.designation
CID1122.employee
Run Code Online (Sandbox Code Playgroud)

我面对卡夫卡 java.lang.NullPointerException

connect_1    | 2019-01-30 06:14:47,302 INFO   ||  Checking MySql dialect for existence of table "CID1122"."employee"   [io.confluent.connect.jdbc.dialect.MySqlDatabaseDialect]
connect_1    | 2019-01-30 06:14:47,303 INFO   ||  Using MySql dialect table "CID1122"."employee" absent   [io.confluent.connect.jdbc.dialect.MySqlDatabaseDialect]
connect_1    | 2019-01-30 …
Run Code Online (Sandbox Code Playgroud)

regex jdbc apache-kafka apache-kafka-connect confluent-platform

0
推荐指数
1
解决办法
3580
查看次数