标签: confluent-platform

Avro.AvroException:数组未在字段中实现非通用 IList

我有 .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

我是否做错了什么或从我的架构中遗漏了什么?

c# avro apache-kafka confluent-platform asp.net-core-3.1

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

在终端中读取来自 Kafka 的 Avro 消息 - kafka-avro-console-consumer 替代方案

我正在尝试找到如何以可读格式从 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 消息?

avro apache-kafka confluent-platform

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

Kafka 连接教程停止工作

我在此链接中遵循步骤 #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

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

Kafka Producer 配置重试策略

需要更改 Kafka Producer 配置的哪些参数,以便生产者应该: 1)重试 n 次 2)在同一消息的 n 间隔后,以防代理宕机。

我需要处理与此相关的情况:https : //github.com/rsyslog/rsyslog/issues/1052

python rsyslog apache-kafka kafka-producer-api confluent-platform

2
推荐指数
2
解决办法
8046
查看次数

模式注册中的 identityMapCapacity 是什么意思

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

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

如何将 confluent-kafka 与密钥存储文件一起使用

当我使用密钥存储文件时,添加属性

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 中找不到该属性

c# ssl truststore apache-kafka confluent-platform

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

Kafka Connect:找不到合适的驱动程序

我正在尝试使用 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

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

用于融合平台的 sbt 解析器

我无法在我的 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)

scala sbt apache-kafka confluent-platform

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

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

如何跨团队共享 avro 模式定义

Kafka schema-registry 提供了一种使用通用数据契约对来自 Kafka 的数据进行序列化和反序列化的好方法。然而,数据契约(.avsc 文件)是生产者和消费者之间的粘合剂。

一旦制作人制作了 .avsc 文件,就可以将其签入制作人一侧的版本控制。根据语言,它也会自动生成类。

然而,

  1. 消费者下拉模式定义以供参考的最佳机制是什么?有没有像 swaggerhub 或 avro 的典型 api 文档门户之类的东西?
  2. 如果我们使用 Confluent 平台,控制中心提供了一个 gui 来查看与主题关联的模式,但它也允许用户进行编辑。生产者和消费者团队之间将如何工作?什么会阻止消费者或任何人直接在 Confluent 平台上编辑模式?
  3. 这是我们需要使用rest-proxy自定义构建的东西吗?

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

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