标签: confluent-platform

在 alpine 容器中使用 confluence-kafka python 客户端

我正在尝试运行一个与 kafka 通信的简单 python 应用程序。我正在寻找一个高山容器。这是我当前的 dockerfile(它不是最佳的......只是想让事情暂时正常工作)。

FROM python:3.6-alpine
MAINTAINER Ashic Mahtab (ashic@live.com)

RUN mkdir -p /usr/src/app
WORKDIR /usr/src/app

RUN echo "http://dl-cdn.alpinelinux.org/alpine/edge/community" >> /etc/apk/repositories
RUN apk update && apk --no-cache add librdkafka


COPY requirements.txt /usr/src/app/
RUN pip install --no-cache-dir -r requirements.txt

COPY api /usr/src/app/api
COPY static /usr/src/app/static

CMD ["python", "api/index.py"]
Run Code Online (Sandbox Code Playgroud)

需求文件中有 confluence-kafka 。构建失败

OK: 8784 distinct packages available
fetch http://dl-cdn.alpinelinux.org/alpine/v3.4/main/x86_64/APKINDEX.tar.gz
fetch http://dl-cdn.alpinelinux.org/alpine/v3.4/community/x86_64/APKINDEX.tar.gz
fetch http://dl-cdn.alpinelinux.org/alpine/edge/community/x86_64/APKINDEX.tar.gz
ERROR: unsatisfiable constraints:
  so:libcrypto.so.41 (missing):
    required by:
                 librdkafka-0.9.4-r1[so:libcrypto.so.41]
                 librdkafka-0.9.4-r1[so:libcrypto.so.41]
                 librdkafka-0.9.4-r1[so:libcrypto.so.41]
  so:libssl.so.43 (missing):
    required by:
                 librdkafka-0.9.4-r1[so:libssl.so.43]
                 librdkafka-0.9.4-r1[so:libssl.so.43]
                 librdkafka-0.9.4-r1[so:libssl.so.43] …
Run Code Online (Sandbox Code Playgroud)

apache-kafka docker kafka-python alpine-linux confluent-platform

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

Kafka REST 代理 API 有哪些好处?

我不知道Kafka REST Proxy API的优点。它是一个 REST API,所以我知道它对于管理来说很方便。人们为什么使用 Kafka REST 代理 API?添加对生产者或消费者的 Maven 依赖是否很麻烦?

另外,我知道kafka客户端有更好的性能。

rest apache-kafka kafka-rest confluent-platform

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

Confluent 的 Kafka REST 代理与 Kafka 客户端

很好奇Confluent的Kafka REST Proxy和用kafka官方客户端库实现的生产者/消费者的优缺点。我知道 Confluent 的 Kafka REST 代理用于管理任务和 kafka 客户端不支持的语言。

那么,kafka客户端有哪些优势呢?

apache-kafka kafka-rest confluent-platform

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

Confluence:错误无法对表 TimestampIncrementingTableQuerier mysql-jdbc 运行查询

我正在尝试对 MySQL 使用模式时间戳,行数有限,因为我的表大小为 2.6 GB。

以下是我正在使用的连接器属性:

{
        "name": "jdbc_source_mysql_registration_query",
        "config": {
                 "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
                 "key.converter": "io.confluent.connect.avro.AvroConverter",
                 "key.converter.schema.registry.url": "http://localhost:8081",
                 "value.converter": "io.confluent.connect.avro.AvroConverter",
                 "value.converter.schema.registry.url": "http://localhost:8081",
                 "connection.url": "jdbc:mysql://localhost:3310/users?zeroDateTimeBehavior=ROUND&useCursorFetch=true&defaultFetchSize=1000&user=kotesh&password=kotesh",
                 "query": "SELECT matriid,DateUpdated  from users.employee WHERE date(DateUpdated)>='2018-11-28' ",
                 "mode": "timestamp",
                 "timestamp.column.name": "DateUpdated",
                 "validate.non.null": "false",
                 "topic.prefix": "mysql-prod-kot-"
        }
}
Run Code Online (Sandbox Code Playgroud)

我得到如下:

INFO TimestampIncrementingTableQuerier{table=null,query='SELECT matriid,DateUpdated from users.employee WHERE date(DateUpdated)>='2018-11-28'',topicPrefix='mysql-prod-kot-',incrementingColumn='', timestampColumns=[DateUpdated]} 准备好的 SQL 查询: SELECT matriid,DateUpdated from users.employee WHERE date(DateUpdated)>='2018-11-28' WHERE DateUpdated> ? 和 DateUpdated< ? ORDER BY DateUpdatedASC (io.confluence.connect.jdbc.source.TimestampIncrementingTableQuerier:161) [2018-11-29 17:29:00,981] 错误无法运行表 TimestampIncrementingTableQuerier{table=null, query='SELECT matriid,DateUpdated 的查询来自 users.employee …

jdbc apache-kafka apache-kafka-connect confluent-schema-registry confluent-platform

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

如何在java中的kafka中获取消费者组的消费者延迟

我想知道使用 java 的消费者组的消费者滞后。我尝试过使用

kafka-consumer-groups --describe --bootstrap-server localhost:9092 --group MyGroupName
Run Code Online (Sandbox Code Playgroud)

并且滞后是可见的。

我如何在 Java 中执行此操作?

我尝试过使用org.apache.kafka.clients.admin.AdminClient,但无法获得每个消费者组的延迟。

我在用 confluent 5.0.1 which has kafka 2.0.1 org.apache.kafka - kafka-clients - 2.0.1

java apache-kafka kafka-consumer-api confluent-platform

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

kafka jdbc接收器连接器中的批量大小

我想通过jdbc接收器批量读取5000条记录,为此我在jdbc接收器配置文件中使用了batch.size:

name=jdbc-sink
connector.class=io.confluent.connect.jdbc.JdbcSinkConnector
tasks.max=1
batch.size=5000
topics=postgres_users

connection.url=jdbc:postgresql://localhost:34771/postgres?user=foo&password=bar
file=test.sink.txt
auto.create=true
Run Code Online (Sandbox Code Playgroud)

但是,batch.size 不起作用,因为当新记录插入源数据库时,记录也会插入数据库。

如何实现批量插入5000个?

connector apache-kafka apache-kafka-connect confluent-platform

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

是否可以从 kafka 消息中获取消息密钥的最新值

假设我对同一个消息键有不同的值。

例如:

{
userid: 1,
email: user123@xyz.com }

{
userid: 1,
email: user456@xyz.com }

{
userid: 1,
email: user789@xyz.com }
Run Code Online (Sandbox Code Playgroud)

在上面的这种情况下,我只想要用户更新的最新值,即“user789@xyz.com”。

我的 kafka 流应该只给我第三个值,而不是前两个值。

apache-kafka apache-kafka-streams spring-kafka confluent-platform ksqldb

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

Cloud Confluent kafka 安装在 ubuntu 18 中失败

说明页面显示为:-

  1. 安装 Confluent Cloud CLI 运行此命令以安装 Confluent Cloud CLI。curl -L --http1.1 https://cnfl.io/ccloud-cli | sh -s -- -b /usr/local/bin

权限失败

“安装:无法创建常规文件‘/usr/local/bin/ccloud’:权限被拒绝”

即使我尝试过

sudo curl -L --http1.1 https://cnfl.io/ccloud-cli | sh -s -- -b /usr/local/bin 
Run Code Online (Sandbox Code Playgroud)

同样的错误。我该如何安装?

confluent-cloud confluent-platform

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

Kafka - max.in.flight.requests.per.connection 是每个生产者或会话?

我正在浏览文档,对参数“max.in.flight.requests.per.connection”有点困惑

客户端在阻塞之前在单个连接上发送的未确认请求的最大数量。请注意,如果此设置设置为大于 1 并且存在发送失败的情况,则存在由于重试而导致消息重新排序的风险(即,如果启用了重试)。

短语“未确认的请求”是指每个生产者或每个连接或每个客户端?

apache-kafka confluent-platform

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

KafkaException:在 kafka.tools.ConsoleProducer$LineMessageReader.readMessage 第 1 行找不到密钥

简而言之,我启动了 Kafka,成功创建了一个主题,启动了一个启用了密钥的生产者。到目前为止,一切都很好。我发送一条简单的消息,然后我得到

root@kafka:/# kafka-console-producer --broker-list localhost:9092 --topic testkey --property "parse.key=true" propert
y "key.separator=:"
>1:fisrtmessage
org.apache.kafka.common.KafkaException: No key found on line 1: 1:fisrtmessage
        at kafka.tools.ConsoleProducer$LineMessageReader.readMessage(ConsoleProducer.scala:275)
        at kafka.tools.ConsoleProducer$.main(ConsoleProducer.scala:55)
        at kafka.tools.ConsoleProducer.main(ConsoleProducer.scala)
root@kafka:/#
Run Code Online (Sandbox Code Playgroud)

这里既是生产者又是消费者。正如你所看到的,消费者没有收到消息,生产者崩溃了。

命令屏幕

docker-compose.yml

version: '3'

services:
  zookeeper:
    image: confluentinc/cp-zookeeper:5.4.0
    hostname: zookeeper
    container_name: zookeeper
    ports:
      - "2181:2181"
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000

  broker:
    image: confluentinc/cp-server:5.4.0
    hostname: broker
    container_name: broker
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181'
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092
      KAFKA_METRIC_REPORTERS: io.confluent.metrics.reporter.ConfluentMetricsReporter
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
      KAFKA_CONFLUENT_LICENSE_TOPIC_REPLICATION_FACTOR: 1
      CONFLUENT_METRICS_REPORTER_BOOTSTRAP_SERVERS: …
Run Code Online (Sandbox Code Playgroud)

apache-kafka docker kafka-producer-api apache-zookeeper confluent-platform

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