标签: confluent-platform

如何从 confluence_python AVRO 消费者获取最新的偏移值

我对 confluence_kafka 还很陌生,但我已经获得了一些使用 kafka-python 的经验。我想做的是改变开始消费消息的偏移量。这就是为什么我想构建一个能够返回到以前的消息的消费者客户端,以便返回将填充仪表板的数据。说使用kafka-python包我可以使用seek_to_endhttps://github.com/dpkp/kafka-python/blob/c0fddbd24269d4333e3b6630a23e86ffe33dfcb6/kafka/consumer/group.py#L788)方法来获取位置值最新的提交。这样我就可以使用该seek方法减去值并返回到之前的消息(https://github.com/dpkp/kafka-python/blob/c0fddbd24269d4333e3b6630a23e86ffe33dfcb6/kafka/consumer/group.py#L738

另一方面,conflient_kafka似乎没有类似的功能,到目前为止我发现的是使用变量OFFSET_END,其值为-1,并且它不会返回最新和最大的偏移数值一。我也可以使用“seek”函数,但我需要一种方法来获取最新偏移量的数值,而不是-1.

我的 avro 消费者看起来像

from confluent_kafka.avro import AvroConsumer

if __name__ == '__main__':
     c = AvroConsumer({"bootstrap.servers": "locahost:29092", "group.id":"mygroup",'schema.registry.url': 'http://localhost:8081',
                  'enable.auto.commit': True,'default.topic.config': {'auto.offset.reset': 'smallest'}})

def my_assign (consumer, partitions):
    for p in partitions:
        p.offset = confluent_kafka.OFFSET_END
        print("offset=",p.offset)
    print('assign', partitions)
    print('position:',consumer.position(partitions))
    consumer.assign(partitions)

c.subscribe(["mytopic"],on_assign=my_assign)

while True:
    m = c.poll(1)
    if m is None:
        continue

    if m.error() is None:
        print('Received message', m.value(),m.offset())
c.close()
Run Code Online (Sandbox Code Playgroud)

产生以下结果:

offset= -1
assign [TopicPartition{topic=mytopic,partition=0,offset=-1,error=None}]
position: [TopicPartition{topic=mytopic,partition=0,offset=-1001,error=None}]
Run Code Online (Sandbox Code Playgroud)

并等待下一条消息。我想知道是否有人可以帮助我。谢谢

python kafka-python confluent-platform

6
推荐指数
1
解决办法
7958
查看次数

Kafka连接HDFS接收器错误无法创建WAL

我正在使用 Kafka 连接 HDFS。当我尝试运行连接器时,出现以下异常:

错误无法创建 WAL 编写器:无法为客户端 [IP] 的 [DFSClient_NONMAPREDUCE_208312334_41] 创建文件 [/path/log],因为该文件已由 [DFSClient_NONMAPREDUCE_165323242_41] 创建

请问有什么建议吗?

hdfs apache-kafka apache-kafka-connect confluent-platform

6
推荐指数
1
解决办法
838
查看次数

如何配置Confluence Platform Kafka连接日志?

我正在使用融合的 kafka 连接服务,但它没有写入日志/var/log/kafka。如何配置它以便它写入日志/var/log/kafka

目前 /var/log/kafka 只有以下日志文​​件 -

-rw-r--r-- 1 cp-kafka confluent     0 Sep 20 14:51 kafka-request.log
-rw-r--r-- 1 cp-kafka confluent     0 Sep 20 14:51 kafka-authorizer.log
-rw-r--r-- 1 cp-kafka confluent  1622 Nov 13 15:43 log-cleaner.log
-rw-r--r-- 1 cp-kafka confluent  7611 Nov 13 20:57 state-change.log
-rw-r--r-- 1 cp-kafka confluent  1227 Nov 14 11:13 server.log
-rw-r--r-- 1 cp-kafka confluent 16683 Nov 14 11:13 controller.log
Run Code Online (Sandbox Code Playgroud)

当进一步检查时,我发现日志写入/var/log/messages(我不想要)。看看下面connect-log4j.properties

log4j.rootLogger=INFO, stdout
log4j.appender.stdout=org.apache.log4j.ConsoleAppender
log4j.appender.stdout.layout=org.apache.log4j.PatternLayout
log4j.appender.stdout.layout.ConversionPattern=[%d] %p %m (%c:%L)%n
log4j.logger.org.apache.zookeeper=ERROR
log4j.logger.org.I0Itec.zkclient=ERROR …
Run Code Online (Sandbox Code Playgroud)

apache-kafka apache-kafka-connect confluent-platform

6
推荐指数
1
解决办法
5951
查看次数

KEYSTORE.JKS 存在失败 - 退出代码为 1 #662 - Confluence kafka

我正在尝试将 ssl 配置为汇合的 kafka docker 平台,并在开始说时出现错误

日志:

命令 [/usr/local/bin/dub 路径 /etc/kafka/secrets/kafka.server.keystore.jks 存在] 失败!kafka_kafka-broker1_1_13d7835ad32d 退出,代码为 1

码头工人配置:

version:  '3'
services:
  zookeeper1:
    image: confluentinc/cp-zookeeper:5.1.0
    hostname: zookeeper1
    ports:
      - "2181:2181"
      - "2888:2888"
      - "3888:3888"
    environment:
      ZOOKEEPER_SERVER_ID: 1
      ZOOKEEPER_SERVERS:  0.0.0.0:2888:3888
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
    logging:  
      driver: "json-file"
      options:
        max-size: "10m"
        max-file: "3"
    volumes:
      - zookeeper-data:/var/lib/zookeeper/data
      - zookeeper-log:/var/lib/zookeeper/log
  kafka-broker1:
    image: confluentinc/cp-kafka:5.1.0
    hostname: kafka-broker1:
    ports:
      - "9092:9092"
      - "9093:9093"
    environment:
      KAFKA_LISTENERS: "PLAINTEXT://0.0.0.0:9092,SSL://0.0.0.0:9093"
      KAFKA_ADVERTISED_LISTENERS: "PLAINTEXT://kafkassl.com:9092,SSL://kafkassl.com:9093"
      KAFKA_ZOOKEEPER_CONNECT: zookeeper1:2181
      KAFKA_LOG4J_LOGGERS: "kafka.controller=INFO,kafka.producer.async.DefaultEventHandler=INFO,state.change.logger=INFO"
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 2
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false"
      KAFKA_DELETE_TOPIC_ENABLE: "true"
      KAFKA_LOG_RETENTION_HOURS: 168 …
Run Code Online (Sandbox Code Playgroud)

ssl jks apache-kafka docker confluent-platform

6
推荐指数
1
解决办法
4149
查看次数

Confluence Replicator 无法重新配置连接器任务?

我过去使用过镜像制作器而不是 Replicator,并且收到错误,但现在确定从哪里开始调试它。

这是错误:

[2019-08-12 18:04:09,672] ERROR Failed to reconfigure connector's 
tasks, retrying after backoff: (org.apache.kafka.connect.runtime.distributed.DistributedHerder:958) 
org.apache.kafka.connect.errors.ConnectException: Could not obtain timely topic metadata update from source cluster
    at io.confluent.connect.replicator.TopicMonitorThreadWithZk.assignments(TopicMonitorThreadWithZk.java:138)
    at io.confluent.connect.replicator.ReplicatorSourceConnector.taskConfigs(ReplicatorSourceConnector.java:99)
    at org.apache.kafka.connect.runtime.Worker.connectorTaskConfigs(Worker.java:317)
    at org.apache.kafka.connect.runtime.distributed.DistributedHerder.reconfigureConnector(DistributedHerder.java:997)
    at org.apache.kafka.connect.runtime.distributed.DistributedHerder.reconfigureConnectorTasksWithRetry(DistributedHerder.java:950)
    at org.apache.kafka.connect.runtime.distributed.DistributedHerder.startConnector(DistributedHerder.java:914)
    at org.apache.kafka.connect.runtime.distributed.DistributedHerder.access$1300(DistributedHerder.java:110)
    at org.apache.kafka.connect.runtime.distributed.DistributedHerder$15.call(DistributedHerder.java:924)
    at org.apache.kafka.connect.runtime.distributed.DistributedHerder$15.call(DistributedHerder.java:920)
    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)

apache-kafka confluent-platform

6
推荐指数
0
解决办法
1237
查看次数

Kafka Connect JDBC Sink 连接器 - java.sql.SQLException:找不到合适的驱动程序

我正在尝试在docker的帮助下使用kafka debezium(Kafka流)将一个数据库的表数据下沉到另一个数据库。数据库流工作正常。但流式数据接收另一个 MySQL DB 进程时出现错误。

对于我的连接器接收器配置如下。

 {
  "name": "mysql_sink",
  "config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
    "topics": "mysql-connect.kafka_test.employee",
    "connection.url": "jdbc:mysql://localhost/kafka_test_1&user=debezium&password=xxxxx",
    "auto.create": "true",
    "auto.evolve": "true",
    "insert.mode": "upsert",
    "pk.fields": "id",
    "pk.mode": "record_value",
    "errors.tolerance": "all",
    "errors.log.enable":"true",
    "errors.log.include.messages":"true",
    "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",
    "name": "mysql_sink"
  }
}
Run Code Online (Sandbox Code Playgroud)

但我收到错误。

org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception.
org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:560)
org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:321)
org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:224)
org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:192)
org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:175)
org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:219)
java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
java.util.concurrent.FutureTask.run(FutureTask.java:266)
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
java.lang.Thread.run(Thread.java:748)\nCaused by: org.apache.kafka.connect.errors.ConnectException: java.sql.SQLException: No suitable driver found for jdbc:mysql://localhost/kafka_test_1&user=debezium&password=xxxxx
io.confluent.connect.jdbc.util.CachedConnectionProvider.getValidConnection(CachedConnectionProvider.java:59)
io.confluent.connect.jdbc.sink.JdbcDbWriter.write(JdbcDbWriter.java:52)
io.confluent.connect.jdbc.sink.JdbcSinkTask.put(JdbcSinkTask.java:66)
org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:538)\n\t... 10 more\nCaused by: java.sql.SQLException: No …
Run Code Online (Sandbox Code Playgroud)

apache-kafka apache-kafka-connect confluent-platform

6
推荐指数
1
解决办法
4171
查看次数

confluence-kafka-python:当 Broker 不可用时如何使初始连接超时?

我正在使用confluent-kafka-python,当我尝试连接到已关闭的代理时,发现它无限挂起。我似乎无法应用在文档中找到的任何超时设置:

from confluent_kafka import Consumer

conf = {'bootstrap.servers': f"{self.host}:{self.port}",
           'group.id': "foo",
           'auto.offset.reset': 'smallest',
           'socket.timeout.ms':'2000', 'socket.max.fails':2,
           'metadata.request.timeout.ms': 5000,
           'reconnect.backoff.max.ms':'5000',
           'api.version.request.timeout.ms':'5000',
           #api.version.fallback.ms
           'session.timeout.ms':'2000',
           #heartbeat.interval.ms
           'coordinator.query.interval.ms':'1000',
           #max.poll.interval.ms
           #auto.commit.interval.ms,
           "debug":"generic, broker, topic, metadata",

   }

try:
     self.consumer = Consumer(conf)
Run Code Online (Sandbox Code Playgroud)

我在日志中得到:

%7|1584702589.065|CONNECT|rdkafka#consumer-1| [thrd:x.x.x.x:6667/bootstrap]: x.x.x.x:6667/bootstrap: broker in state TRY_CONNECT connecting
%7|1584702589.065|STATE|rdkafka#consumer-1| [thrd:x.x.x.x:6667/bootstrap]: x.x.x.x:6667/bootstrap: Broker changed state TRY_CONNECT -> CONNECT
%7|1584702589.065|BROADCAST|rdkafka#consumer-1| [thrd:x.x.x.x:6667/bootstrap]: Broadcasting state change
%7|1584702589.065|CONNECT|rdkafka#consumer-1| [thrd:x.x.x.x:6667/bootstrap]: x.x.x.x:6667/bootstrap: Connecting to ipv4#x.x.x.x:6667 (plaintext) with socket 11
%7|1584702589.065|CONNECT|rdkafka#consumer-1| [thrd:app]: Cluster connection already in progress: application metadata request
%7|1584702589.066|CONNECT|rdkafka#consumer-1| …
Run Code Online (Sandbox Code Playgroud)

apache-kafka kafka-consumer-api kafka-python confluent-platform

6
推荐指数
0
解决办法
3770
查看次数

Kafka - min.insync.replicas 解释

我正在浏览文档,查看多个地方,这增加了混乱..

关于属性 min.insync.replicas

当生产者将 acks 设置为“全部”(或“-1”)时,此配置指定必须确认写入才能将写入视为成功的最小副本数。如果无法满足此最小值,则生产者将引发异常(NotEnoughReplicas 或 NotEnoughReplicasAfterAppend)。当一起使用时,min.insync.replicas 和 acks 允许您强制执行更大的持久性保证。一个典型的场景是创建一个复制因子为 3 的主题,将 min.insync.replicas 设置为 2,并使用“all”的 acks 进行生产。如果大多数副本没有收到写入,这将确保生产者引发异常。

我提出的问题,

  1. 此属性仅在作为“发送记录”(生产者)的一部分与“确认”一起使用时才有意义,还是作为消费者流的一部分也有任何影响?
  2. 如果 acks=all 和 min.insync.replicas = 1(default value :1 ) --> 与 acks = 1 一样吗?(考虑复制因子 3 ?

更新 #1 我遇到了这个短语

“当生产者指定 ack (-1 / all config) 时,它仍然会等待当时所有同步副本的 ack(独立于最小同步副本的设置)。因此,如果您在 4 个副本同步时发布那么除非所有 4 个副本都提交消息(即使最小同步副本配置为 2),否则您将不会收到确认。”

这句话是如何与今天相关的?这个属性“最小同步副本”是否仍然独立?

apache-kafka confluent-platform

6
推荐指数
1
解决办法
3625
查看次数

融合模式注册表:使用具有单个属性的对象发布简单的 JSON 模式

OS:  Ubuntu 18.x
docker image (from dockerhub.com, as of 2020-09-25):  confluentinc/cp-schema-registry:latest
Run Code Online (Sandbox Code Playgroud)

我正在探索 Confluence 架构注册表的 HTTP API。首先,是否存在关于注册表假定的 JSON 架构定义版本的明确断言?目前,我假设Draft v7.0。更广泛地说,我认为返回受支持架构的 API 应该列出版本。例如,而不是:

$ curl -X GET http://localhost:8081/schemas/types
["JSON","PROTOBUF","AVRO"]
Run Code Online (Sandbox Code Playgroud)

你将会拥有:

$ curl -X GET http://localhost:8081/schemas/types
[{"flavor": "JSON", "version": "7.0"}, {"flavor": "PROTOBUF", "version": "1.2"}, {"flavor": "AVRO", "version": "3.5"}]
Run Code Online (Sandbox Code Playgroud)

所以至少程序员会确切地知道模式注册表的假设。

抛开这个问题不谈,我似乎无法将一个相当简单的 JSON 模式发布到注册表:

$ curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" --data '{ "schema": "{ \"type\": \"object\", \"properties\": { \"f1\": { \"type\": \"string\" } } }" }' http://localhost:8081/subjects/mytest-value/versions
{"error_code":42201,"message":"Either the input schema or …
Run Code Online (Sandbox Code Playgroud)

confluent-schema-registry confluent-platform

6
推荐指数
1
解决办法
3387
查看次数

使用 Python 获取 Confluence Kafka 主题的最新消息

到目前为止,这是我尝试过的:

from confluent_kafka import Consumer

c = Consumer({... several security/server settings skipped...
              'auto.offset.reset': 'beginning',
              'group.id': 'my-group'})

c.subscribe(['my.topic'])
msg = poll(30.0)  # msg is of None type.
Run Code Online (Sandbox Code Playgroud)

msg几乎总是最终成为None这样。我认为问题可能是'my-group'已经消耗了所有消息'my.topic'......但我不在乎消息是否已经被消耗 - 我仍然需要最新的消息。具体来说,我需要最新消息的时间戳。

我又尝试了一些,从这里看来,该主题中可能有 25 条消息,但我不知道如何获取它们:

a = c.assignment()
print(a)  # Outputs [TopicPartition{topic=my.topic,partition=0,offset=-1001,error=None}]
offsets = c.get_watermark_offsets(a[0])
print(offsets)  # Outputs: (25, 25)
Run Code Online (Sandbox Code Playgroud)

如果因为该主题从未写入任何内容而没有消息,我该如何确定?如果是这样,我如何确定该主题存在了多长时间?我正在编写一个脚本,自动删除过去 X 天内未写入的任何主题(最初为 14 个 - 可能会随着时间的推移进行调整。)

python apache-kafka kafka-consumer-api confluent-platform confluent-kafka-python

6
推荐指数
1
解决办法
6289
查看次数