我正在使用以下水槽。问题是它将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
我正在尝试使用 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 还很陌生,请帮忙!
我已经使用以下命令启动了 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)
我怎样才能找到要使用哪个端口?
嗨,我正在使用 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
在 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) 我有以下设置 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
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) 我的应用程序控制台日志显示了以下很多内容:
%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) 我已经使用 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) 我有开始使用 Kafka 的用例,并且正在寻找开源免费(生产)kafka。
当检查Confluence 7.1平台看起来很合适时,因为它捆绑了zookeeper/kafka/schema注册表/kafka UI。
在决定继续之前,只想检查一下 Confluence Platform 7.1 是否免费且开源?我需要购买许可或付费支持吗?
我正在使用这个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