我正在尝试运行一个与 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
我不知道Kafka REST Proxy API的优点。它是一个 REST API,所以我知道它对于管理来说很方便。人们为什么使用 Kafka REST 代理 API?添加对生产者或消费者的 Maven 依赖是否很麻烦?
另外,我知道kafka客户端有更好的性能。
很好奇Confluent的Kafka REST Proxy和用kafka官方客户端库实现的生产者/消费者的优缺点。我知道 Confluent 的 Kafka REST 代理用于管理任务和 kafka 客户端不支持的语言。
那么,kafka客户端有哪些优势呢?
我正在尝试对 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 BYDateUpdatedASC (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
我想知道使用 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
我想通过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
假设我对同一个消息键有不同的值。
例如:
{
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
说明页面显示为:-
权限失败
“安装:无法创建常规文件‘/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)
同样的错误。我该如何安装?
我正在浏览文档,对参数“max.in.flight.requests.per.connection”有点困惑
客户端在阻塞之前在单个连接上发送的未确认请求的最大数量。请注意,如果此设置设置为大于 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
apache-kafka ×9
docker ×2
kafka-rest ×2
alpine-linux ×1
connector ×1
java ×1
jdbc ×1
kafka-python ×1
ksqldb ×1
rest ×1
spring-kafka ×1