是否有生产者和消费者组在 Java 中使用 Kafka Rest api 的好例子?我不是在寻找生产者和消费者的简单消费者或卡夫卡客户端示例。任何帮助表示赞赏。
apache-kafka kafka-consumer-api kafka-producer-api confluent-platform
我在一组 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
我正在尝试使用 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
我正在尝试使用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) 我有以下 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) 当我启动 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) 目前,我们的经纪商使用 Kafka 0.8.2。我们使用 .Poll() 方法来抓取消息并在收集到 1000 条消息后提交。我们可以使用这个库很好地从集群中消费,并且我可以看到偏移量通过手动提交调用存储在 __consumer_offsets 主题中。但是,当消费者重新启动时,它不会使用存储的偏移量。相反,它从最早的偏移量重新启动(我有 auto.offset.reset=earliest)。我是否需要做一些特定的事情来强制 Kafka 使用这些偏移量而不是在 Zookeeper 中查找,或者消费者应该处理这个问题?有没有办法强制消费者将偏移量写入 Zookeeper 而不是 __consumer_offsets 主题?
连接独立.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) 我是 Confluence/Kafka 的新手,我想从 kafka 中查找元数据信息
我想知道
Confluence版本是5.0
可以提供此信息的类(方法)是什么?
是否有任何 Rest API 用于相同的情况?
还需要连接 Zookeeper 才能获取此信息。
apache-kafka kafka-producer-api confluent-schema-registry confluent-platform
我知道,在 KSQL 中我们可以将偏移量设置为最早或最晚但是我们可以获取特定时间段的数据,即我需要从 2020 年 5 月 6 日起将数据插入到主题中?