我有 .netcore 3.1 api 我正在使用 kafka 和我正在使用的融合云库,它们的版本是:
架构注册表 1.7.0 Serdes 1.3.0
在我的架构中,我有:
{
"name": "orderLineIds",
"type": {
"type": "array",
"items": "int"
},
"default": [],
"doc": "Associated order line item ids related to this promotion"
}
Run Code Online (Sandbox Code Playgroud)
这会在我的班级中生成以下内容:
public IList<System.Int32> orderLineIds
{
get
{
return this._orderLineIds;
}
set
{
this._orderLineIds = value;
}
}
Run Code Online (Sandbox Code Playgroud)
当我尝试产品活动时,出现以下错误:
Avro.AvroException:数组未在字段 orderLineIds 中实现非通用 IList
我是否做错了什么或从我的架构中遗漏了什么?
我正在尝试找到如何以可读格式从 Kafka 主题读取 Avro 消息的最简单方法。kafka-avro-console-consumer可以选择通过以下方式使用 Confluence
./kafka-avro-console-consumer \
--topic topic \
--from-beginning \
--bootstrap-server bootstrap_server_url \
--max-messages 10 \
--property schema.registry.url=schema_registry_url
Run Code Online (Sandbox Code Playgroud)
但为此我需要下载整个 Confluence 平台 (1.7 GB),我认为这在我的场景中是一种过度杀伤力。
有没有其他方法可以轻松地从终端中的 Kafka 主题获取 Avro 消息?
我在此链接中遵循步骤 #7(使用 Kafka Connect 导入/导出数据):
http://kafka.apache.org/documentation.html#quickstart
它运行良好,直到我删除了“test.txt”文件。主要是因为这就是 log4j 文件的工作方式。一段时间后,文件将被旋转 - 我的意思是 - 它将被重命名,并且将开始写入具有相同名称的新文件。
但是之后,我删除了“test.txt”,连接器停止工作。我重新启动了连接器、代理、zookeeper 等,但是“test.txt”中的新行不会进入“connect-test”主题,因此不会进入“test.sink.txt”文件。
我怎样才能解决这个问题?
apache-kafka kafka-consumer-api kafka-producer-api apache-kafka-connect confluent-platform
需要更改 Kafka Producer 配置的哪些参数,以便生产者应该: 1)重试 n 次 2)在同一消息的 n 间隔后,以防代理宕机。
我需要处理与此相关的情况:https : //github.com/rsyslog/rsyslog/issues/1052
python rsyslog apache-kafka kafka-producer-api confluent-platform
identityMapCapacityConfluent Schema Registry中的含义是什么CachedSchemaRegistryClient。根据文档,其声明如下:
public CachedSchemaRegistryClient(@NotNull String baseUrl,int identityMapCapacity)
Run Code Online (Sandbox Code Playgroud)
我看到几个帖子,它用int10初始化,在某个地方它是 1000。所以我不确定它到底是什么意思,我应该使用什么。
java avro apache-kafka confluent-schema-registry confluent-platform
当我使用密钥存储文件时,添加属性
ssl.keystore.location
ssl.keystore.password
ssl.key.password
ssl.truststore.location
ssl.truststore.password
Run Code Online (Sandbox Code Playgroud)
在配置中,它抛出这个错误:
找不到属性 ssl.truststore.location”
如何将 librdkafka 与密钥存储文件一起使用?这让我很烦恼,有人知道如何将 confluent-kafka 与密钥存储文件一起使用吗?
汇合卡夫卡:https : //github.com/confluentinc/confluent-kafka-dotnet/
按照CONFIGURATION.md:https://github.com/edenhill/librdkafka/blob/master/CONFIGURATION.md
我在 CONFIGURATION.md 中找不到该属性
我正在尝试使用 JDBC-sink 连接器将 Kafka 与 Postgres Sink 一起使用。
例外:
INFO Unable to connect to database on attempt 1/3. Will retry in 10000 ms. (io.confluent.connect.jdbc.util.CachedConnectionProvider:91)
java.sql.SQLException: No suitable driver found for jdbc:postgresql://localhost:5432/casb
at java.sql.DriverManager.getConnection(DriverManager.java:689)
at java.sql.DriverManager.getConnection(DriverManager.java:247)
at io.confluent.connect.jdbc.util.CachedConnectionProvider.newConnection(CachedConnectionProvider.java:85)
at io.confluent.connect.jdbc.util.CachedConnectionProvider.getValidConnection(CachedConnectionProvider.java:68)
at io.confluent.connect.jdbc.sink.JdbcDbWriter.write(JdbcDbWriter.java:56)
at io.confluent.connect.jdbc.sink.JdbcSinkTask.put(JdbcSinkTask.java:69)
at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:495)
at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:288)
at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:198)
at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:166)
at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:170)
at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:214)
at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
at java.util.concurrent.FutureTask.run(FutureTask.java:266)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)
Run Code Online (Sandbox Code Playgroud)
Sink.properties:
name=test-sink
connector.class=io.confluent.connect.jdbc.JdbcSinkConnector
tasks.max=1
topics=fp_test
connection.url=jdbc:postgresql://localhost:5432/casb
connection.user=admin
connection.password=***
auto.create=true
Run Code Online (Sandbox Code Playgroud)
我已经设定 plugin.path=/usr/share/java/kafka-connect-jdbc
在/usr/share/java/kafka-connect-jdbc我有以下文件:
kafka-connect-jdbc-4.0.0.jar …
postgresql apache-kafka apache-kafka-connect confluent-platform
我无法在我的 sbt 中添加融合的 repo。我查看了pom 示例 并找到了在 maven 中添加 repo 的定义。
<repositories>
<repository>
<id>confluent</id>
<url>https://packages.confluent.io/maven/</url>
</repository>
<!-- further repository entries here -->
</repositories>
Run Code Online (Sandbox Code Playgroud)
和依赖
<dependencies>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka_2.11</artifactId>
<version>2.0.0-cp1</version>
</dependency>
<!-- further dependency entries here -->
</dependencies>
Run Code Online (Sandbox Code Playgroud)
我用了
resolvers += Resolver.url("confluent", url("http://packages.confluent.io/maven/")) in build.sbt`
Run Code Online (Sandbox Code Playgroud)
并将依赖项声明为
libraryDependencies += "org.apache.kafka" % "kafka-clients" % "2.0.0-cp1"
libraryDependencies += "org.apache.kafka" %% "kafka" % "2.0.0-cp1"
Run Code Online (Sandbox Code Playgroud)
我仍然得到
::::::::::::::::::::::::::::::::::::::::::::::
[warn] :: UNRESOLVED DEPENDENCIES ::
[warn] ::::::::::::::::::::::::::::::::::::::::::::::
[warn] :: org.apache.kafka#kafka-clients;2.0.0-cp1: not found
[warn] :: org.apache.kafka#kafka_2.12;2.0.0-cp1: not found
[warn] :::::::::::::::::::::::::::::::::::::::::::::: …Run Code Online (Sandbox Code Playgroud) 我想知道 HDF 套件中嵌入的 kafka 和 Confluent 套件的区别,特别是模式注册工具。
apache-kafka hortonworks-data-platform confluent-schema-registry confluent-platform
Kafka schema-registry 提供了一种使用通用数据契约对来自 Kafka 的数据进行序列化和反序列化的好方法。然而,数据契约(.avsc 文件)是生产者和消费者之间的粘合剂。
一旦制作人制作了 .avsc 文件,就可以将其签入制作人一侧的版本控制。根据语言,它也会自动生成类。
然而,
avro apache-kafka confluent-schema-registry confluent-platform
apache-kafka ×10
avro ×4
c# ×2
java ×1
postgresql ×1
python ×1
rsyslog ×1
sbt ×1
scala ×1
ssl ×1
truststore ×1