标签: confluent-platform

卡夫卡休息的例子

是否有生产者和消费者组在 Java 中使用 Kafka Rest api 的好例子?我不是在寻找生产者和消费者的简单消费者或卡夫卡客户端示例。任何帮助表示赞赏。

apache-kafka kafka-consumer-api kafka-producer-api confluent-platform

5
推荐指数
1
解决办法
2万
查看次数

Kafka接收器连接器:即使重新启动后也没有分配任务

我在一组 Docker 容器中使用 Confluence 3.2,其中一个容器正在运行 kafka-connect 工作线程。

由于我尚不清楚的原因,我的四个连接器中的两个 - 具体来说,hpgraphsl 的MongoDB 接收器连接器- 停止工作。我能够确定主要问题:连接器没有分配任何任务,通过调用 可以看出GET /connectors/{my_connector}/status。其他两个连接器(同一类型)没有受到影响,并且正在愉快地产生输出。

我尝试了三种不同的方法来通过 REST API 让连接器再次运行:

  • 暂停和恢复连接器
  • 重新启动连接器
  • 使用相同的配置删除并创建同名的连接器

所有方法都不起作用。我终于让我的连接器再次工作:

  • 删除连接器并以不同的名称创建连接器,my_connector_v2例如my_connector

这里发生了什么?为什么我无法重新启动现有连接器并让它启动实际任务?kafka-connect 工作线程或 Kafka 代理的某些与 kafka-connect 相关的主题中是否有任何陈旧数据需要清理?

我已经在特定连接器的 github 存储库上提交了一个问题,但我觉得这实际上可能是与 kafka-connect 的内在相关的一般错误。有任何想法吗?

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

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

Kafka-ES-Sink : ConnectException: Key 用作文档 id 并且不能为 null

我正在尝试使用 SMT 函数添加密钥以将其用作 ES 文档的文档 ID,但它不起作用。我正在使用 confluence es 连接器。配置文件如下: connect-standalone.properties

bootstrap.servers=localhost:9092
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=false
value.converter.schemas.enable=false
internal.key.converter=org.apache.kafka.connect.json.JsonConverter
internal.value.converter=org.apache.kafka.connect.json.JsonConverter
internal.key.converter.schemas.enable=false
internal.value.converter.schemas.enable=false
offset.storage.file.filename=/tmp/connect.offsets
offset.flush.interval.ms=10000
Run Code Online (Sandbox Code Playgroud)

连接器配置:

#Connector name                                                                  
name=logs-=false                                                                 
#Connector class                                                                 
connector.class=io.confluent.connect.elasticsearch.ElasticsearchSinkConnector    
tasks.max=1                                                                      
topics=test                                                                  
topic.index.map=test:activity                                     
connection.url=http://localhost:9200                                             
type.name=Activity                                                      
#ignore key and schema                                                           
key.ignore=false                                                                 
schema.ignore=true                                                               
transforms=InsertKey,ExtractId                                                   
transforms.InsertKey.type=org.apache.kafka.connect.transforms.ValueToKey         
transforms.InsertKey.fields=recordId                                          
transforms.ExtractId.type=org.apache.kafka.connect.transforms.ExtractField$Key   
transforms.ExtractId.field=recordId 
Run Code Online (Sandbox Code Playgroud)

我正在向 kafka 发送以下消息:

{"recordId":"999","activity":"test","description":"test Cont"}
Run Code Online (Sandbox Code Playgroud)

在接收器连接器中出现此错误:

org.apache.kafka.connect.errors.ConnectException: Key is used as document id and can not be null.
        at io.confluent.connect.elasticsearch.DataConverter.convertKey(DataConverter.java:56)
        at io.confluent.connect.elasticsearch.DataConverter.convertRecord(DataConverter.java:86)
        at io.confluent.connect.elasticsearch.ElasticsearchWriter.write(ElasticsearchWriter.java:210)
        at io.confluent.connect.elasticsearch.ElasticsearchSinkTask.put(ElasticsearchSinkTask.java:119)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:384)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:240)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:172)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:143)
        at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:140) …
Run Code Online (Sandbox Code Playgroud)

elasticsearch apache-kafka apache-kafka-connect confluent-platform

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

连接器配置不包含连接器类型

我正在尝试使用JDBC 连接器连接到集群上的 PostgreSQL 数据库(该数据库不是由集群直接管理)。

我一直在使用以下命令调用 Kafka Connect:

connect-standalone.sh worker.properties jdbc-connector.properties
Run Code Online (Sandbox Code Playgroud)

这是文件的内容worker.properties

class=io.confluent.connect.jdbc.JdbcSourceConnector
name=test-postgres-1
tasks.max=1

internal.key.converter=org.apache.kafka.connect.json.JsonConverter
internal.value.converter=org.apache.kafka.connect.json.JsonConverter
internal.key.converter.schemas.enable=false
internal.value.converter.schemas.enable=false

offset.storage.file.filename=/home/user/offest
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter=org.apache.kafka.connect.json.JsonConverter

connection.url=jdbc:postgresql://database-server.url:port/database?user=user&password=password
Run Code Online (Sandbox Code Playgroud)

这是以下内容jdbc-connector.properties

mode=incrementing
incrementing.column.name=id
topic.prefix=test-postgres-jdbc-
Run Code Online (Sandbox Code Playgroud)

当我尝试使用上述命令启动连接器时,它崩溃并出现以下错误:

[2018-04-16 11:39:08,164] ERROR Failed to create job for jdbc.properties (org.apache.kafka.connect.cli.ConnectStandalone:88)
[2018-04-16 11:39:08,166] ERROR Stopping after connector error (org.apache.kafka.connect.cli.ConnectStandalone:99)
java.util.concurrent.ExecutionException: org.apache.kafka.connect.runtime.rest.errors.BadRequestException: Connector config {mode=incrementing, incrementing.column.name=pdv, topic.prefix=test-postgres-jdbc-} contains no connector type
    at org.apache.kafka.connect.util.ConvertingFutureCallback.result(ConvertingFutureCallback.java:80)
    at org.apache.kafka.connect.util.ConvertingFutureCallback.get(ConvertingFutureCallback.java:67)
    at org.apache.kafka.connect.cli.ConnectStandalone.main(ConnectStandalone.java:96)
Caused by: org.apache.kafka.connect.runtime.rest.errors.BadRequestException: Connector config {mode=incrementing, incrementing.column.name=id, topic.prefix=test-postgres-jdbc-} contains no connector type …
Run Code Online (Sandbox Code Playgroud)

apache-kafka apache-kafka-connect confluent-platform

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

为什么这个测试要花这么长时间?

我有以下 xUnit(版本 2.3.1)测试,用于向 Kafka 发送 100 条消息

[Fact]
public void Test1()
{
    Stopwatch sw = new Stopwatch();

    sw.Start();

    var config = new Dictionary<string, object>
                    {
                        { "group.id", "gid" },
                        { "bootstrap.servers", "localhost" },
                        { "enable.auto.commit", true },
                        { "default.topic.config", new Dictionary<string, object>()
                            {
                                { "message.timeout.ms", 500 }
                            }
                        },
                    };

    var connection = new KafkaConnection(config, new ShreddingQueueCache());

    for (byte i = 0; i < 100; i++)
    {
        connection.Send(new Message(new Guid(1, 2, 3, new byte[] { 0, 1, 0, 1, 0, …
Run Code Online (Sandbox Code Playgroud)

c# xunit .net-core kafka-producer-api confluent-platform

5
推荐指数
0
解决办法
217
查看次数

如何避免有关 ReflectionsException 的 WARN 消息

当我启动 Confluence 时,如何避免在日志中显示 WARN 消息(不将 log4j 级别设置为 ERROR)?plugin.path我已经在属性文件中设置了变量的值${CONFLUENT_HOME}/share/java/kafka-connect-jdbc(最后带有逗号)。我尝试将 kafka-connect-jdbc 存储库放入类路径中,但没有成功。

以下只是日志文件的一小部分示例:

[2018-07-10 15:40:30,168] INFO Reflections took 1 ms to scan 1 urls, producing 5 keys and 6 values [using 1 cores] (org.reflection
s.Reflections)
[2018-07-10 15:40:30,170] WARN could not get type for name org.jmock.Mockery from any class loader (org.reflections.Reflections)
org.reflections.ReflectionsException: could not get type for name org.jmock.Mockery
        at org.reflections.ReflectionUtils.forName(ReflectionUtils.java:390)
        at org.reflections.Reflections.expandSuperTypes(Reflections.java:381)
        at org.reflections.Reflections.<init>(Reflections.java:126)
        at org.apache.kafka.connect.runtime.isolation.DelegatingClassLoader$InternalReflections.<init>(DelegatingClassLoader.java:
365)
        at org.apache.kafka.connect.runtime.isolation.DelegatingClassLoader.scanPluginPath(DelegatingClassLoader.java:277)
        at org.apache.kafka.connect.runtime.isolation.DelegatingClassLoader.scanUrlsAndAddPlugins(DelegatingClassLoader.java:216)
        at org.apache.kafka.connect.runtime.isolation.DelegatingClassLoader.registerPlugin(DelegatingClassLoader.java:208)
        at org.apache.kafka.connect.runtime.isolation.DelegatingClassLoader.initPluginLoader(DelegatingClassLoader.java:177)
        at org.apache.kafka.connect.runtime.isolation.DelegatingClassLoader.initLoaders(DelegatingClassLoader.java:154)
        at …
Run Code Online (Sandbox Code Playgroud)

warnings log4j exception confluent-platform

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

合流的 go kafka 库在重新启动时从最早的偏移量开始

目前,我们的经纪商使用 Kafka 0.8.2。我们使用 .Poll() 方法来抓取消息并在收集到 1000 条消息后提交。我们可以使用这个库很好地从集群中消费,并且我可以看到偏移量通过手动提交调用存储在 __consumer_offsets 主题中。但是,当消费者重新启动时,它不会使用存储的偏移量。相反,它从最早的偏移量重新启动(我有 auto.offset.reset=earliest)。我是否需要做一些特定的事情来强制 Kafka 使用这些偏移量而不是在 Zookeeper 中查找,或者消费者应该处理这个问题?有没有办法强制消费者将偏移量写入 Zookeeper 而不是 __consumer_offsets 主题?

go apache-kafka confluent-platform

5
推荐指数
0
解决办法
427
查看次数

在 Kafka Connect 中找不到 jdbc:mysql 合适的驱动程序

连接独立.properties

connector.class=io.confluent.connect.jdbc.JdbcSourceConnector
bootstrap.servers=10.33.62.20:9092,10.33.62.110:9092,10.33.62.200:9092
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=true
value.converter.schemas.enable=true

offset.storage.file.filename=/tmp/connect.offsets
offset.flush.interval.ms=10000
plugin.path=/grid/1/mukul/confluent-5.0.0/share/java
Run Code Online (Sandbox Code Playgroud)

源-sqlite.properties

name=test-source-sqlite-jdbc-autoincrement
connector.class=io.confluent.connect.jdbc.JdbcSourceConnector
tasks.max=5

connection.url=jdbc:mysql://10.32.177.178:3306/test&user=xxxx&password=xxxxx

table.whitelist=banner_hourly_statistics_v2

group.id=test-mysql-kafka
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter

config.storage.topic=demo-1-distributed-config
offset.storage.topic=demo-1-distributed-offset
status.storage.topic=demo-1-distributed-status

bootstrap.servers=10.33.62.20:9092,10.33.62.110:9092,10.33.62.200:9092
mode=bulk
#incrementing.column.name=id
topic.prefix=test-sqlite-jdbc-
Run Code Online (Sandbox Code Playgroud)

命令:connect-standalone /grid/1/mukul/confluent-5.0.0/etc/kafka/connect-standalone.properties /grid/1/mukul/confluent-5.0.0/etc/kafka-connect-jdbc/source-quickstart-sqlite.properties

在启动日志中,它清楚地显示了加载 JDBC 连接器:

[2018-08-09 06:59:30,072] INFO Loading plugin from: /grid/1/mukul/confluent-5.0.0/share/java/kafka-connect-jdbc (org.apache.kafka.connect.runtime.isolation.DelegatingClassLoader:218)
[2018-08-09 06:59:30,133] INFO Registered loader: PluginClassLoader{pluginLocation=file:/grid/1/mukul/confluent-5.0.0/share/java/kafka-connect-jdbc/} (org.apache.kafka.connect.runtime.isolation.DelegatingClassLoader:241)
[2018-08-09 06:59:30,133] INFO Added plugin 'io.confluent.connect.jdbc.JdbcSinkConnector' (org.apache.kafka.connect.runtime.isolation.DelegatingClassLoader:170)
[2018-08-09 06:59:30,133] INFO Added plugin 'io.confluent.connect.jdbc.JdbcSourceConnector' (org.apache.kafka.connect.runtime.isolation.DelegatingClassLoader:170)
Run Code Online (Sandbox Code Playgroud)

但它失败并出现以下异常:

Invalid value java.sql.SQLException: No suitable driver found for jdbc:mysql://10.32.177.178:3306/test&user=xxxx&password=xxxx for configuration Couldn't open connection to jdbc:mysql://10.32.177.178:3306/test&user=xxxx&password=xxx
Invalid …
Run Code Online (Sandbox Code Playgroud)

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

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

来自kafka的元数据信息

我是 Confluence/Kafka 的新手,我想从 kafka 中查找元数据信息

我想知道

  1. 生产者名单
  2. 主题列表
  3. 主题的架构信息

Confluence版本是5.0

可以提供此信息的类(方法)是什么?
是否有任何 Rest API 用于相同的情况?
还需要连接 Zookeeper 才能获取此信息。

apache-kafka kafka-producer-api confluent-schema-registry confluent-platform

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

我们是否可以选择从特定时间段/时间戳获取 KSQL 流中的数据

我知道,在 KSQL 中我们可以将偏移量设置为最早或最晚但是我们可以获取特定时间段的数据,即我需要从 2020 年 5 月 6 日起将数据插入到主题中?

apache-kafka confluent-platform ksqldb

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