标签: confluent-platform

Docker 中的 Kafka 日志目录

当我在没有 Docker 的情况下运行 kafka 和 Zookeeper 时,我可以在 /tmp/kafka-logs 目录中看到主题分区日志文件。现在使用 Docker,即使我在 docker-compose.yml 的 Volumes 部分中指定了日志目录,我也看不到 docker VM 中的文件,例如“TOPICNAME-PARTITIONNUMBER”..这里有什么我遗漏的吗?知道在 Docker 虚拟机中哪里可以找到这些目录吗?

zookeeper:
 image: confluent/zookeeper
 container_name: zookeeper
 ports:
   - "2181:2181"
   - "15001:15000"
 environment:
 ZK_SERVER_ID: 1
 volumes:
   - /tmp/docker/zk1/logs:/logs
   - /tmp/docker/zk1/data:/data

 kafka1:
  image: confluent/kafka
  container_name: kafka1
  ports:
    - "9092:9092"
    - "15002:15000"
  links:
    - zookeeper
  environment:
    KAFKA_BROKER_ID: 1
    KAFKA_OFFSETS_STORAGE: kafka
    # This is Container IP
    KAFKA_ADVERTISED_HOST_NAME: 192.168.99.100 

  volumes:
    - /tmp/docker/kafka1/logs:/logs
    - /tmp/docker/kafka1/data:/data
Run Code Online (Sandbox Code Playgroud)

apache-kafka docker docker-compose apache-zookeeper confluent-platform

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

将 Avro 架构保存到 Confluence Schema-Registry

我正在尝试使用 Postman 向注册表写入一个非常简单的模式,并且一直很难让它注册。只是注册一个简单的模式真的这么复杂吗?这只是整个过程的第一步,还是我在这里遗漏了一些东西?我正在使用的架构如下:

    {
"schema":{
  "type" : "record",
  "name" : "User",
  "namespace" : "com.temp.avro.model",
  "fields" : [ {
    "name" : "_id",
    "type" : "string"
  }, {
    "name" : "updatedDate",
    "type":"long",
    "logicalType":"timestamp-millis"
  }, {
    "name" : "createdDate",
    "type":"long",
    "logicalType":"timestamp-millis"
  }, {
    "name" : "applicationId",
    "type": ["null", "string"],
    "default": null
  },{
    "name" : "country",
    "type" : "string"
  }, {
    "name" : "bank",
    "type" : "string"
  }]
}
}
Run Code Online (Sandbox Code Playgroud)

我收到以下错误:

Internal Server Error com.fasterxml.jackson.databind.JsonMappingException: Can not deserialize instance of java.lang.String out of START_OBJECT …
Run Code Online (Sandbox Code Playgroud)

avro confluent-schema-registry confluent-platform

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

Kafka-HDFS-Connector- 将数据从 Kafka 发送到 Hadoop

在机器A 中kafka 正在运行,在机器B中安装了 hadoop。现在我想从 kafka 向 hadoop 写入数据。我已经在 Machine A 中安装了 Confluent Platform 。

任何人都可以指导我必须添加哪些配置才能将数据从 kafka 写入在不同机器上运行的 hadoop

hadoop hdfs apache-kafka confluent-platform

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

AVRO 原始类型的 Serde 类

我正在用 Java 编写一个 Kafka 流应用程序,它接受由连接器创建的输入主题,该连接器使用模式注册表和 avro 作为键和值转换器。连接器产生以下模式:

key-schema: "int"
value-schema:{
"type": "record",
"name": "User",
"fields": [
    {"name": "firstname", "type": "string"},
    {"name": "lastname",  "type": "string"}
]}
Run Code Online (Sandbox Code Playgroud)

实际上,有几个主题,key-schema 总是“int”,value-schema 总是某种记录(用户、产品等)。我的代码包含以下定义

Map<String, String> serdeConfig = Collections.singletonMap("schema.registry.url", schemaRegistryUrl);

Serde<User> userSerde = new SpecificAvroSerde<>();
userSerde.configure(serdeConfig, false);
Run Code Online (Sandbox Code Playgroud)

起初我尝试使用类似的东西来消费这个主题, Consumed.with(Serdes.Integer(), userSerde);但这不起作用,因为 Serdes.Integer() 期望使用 4 个字节对整数进行编码,但 avro 使用可变长度编码。使用Consumed.with(Serdes.Bytes(), userSerde);有效,但我真的想要 int 而不是字节,所以我将代码更改为此

KafkaAvroDeserializer keyDeserializer = new KafkaAvroDeserializer()
KafkaAvroSerializer keySerializer = new KafkaAvroSerializer();
keyDeserializer.configure(serdeConfig, true); 
keySerializer.configure(serdeConfig, true);
Serde<Integer> keySerde = (Serde<Integer>)(Serde)Serdes.serdeFrom(keySerializer, keyDeserializer);
Run Code Online (Sandbox Code Playgroud)

这使编译器产生警告(它不喜欢(Serde<Integer>)(Serde)强制转换)但它允许我使用

Consumed.with(keySerde, userSerde); …

java avro apache-kafka apache-kafka-streams confluent-platform

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

在 pkg-config 搜索路径中找不到包 rdkafka

在 pkg-config 搜索路径中找不到包 rdkafka。

Confluent go 包像这样抛出错误

# pkg-config --cflags  -- rdkafka
Package rdkafka was not found in the pkg-config search path.
Perhaps you should add the directory containing `rdkafka.pc'
to the PKG_CONFIG_PATH environment variable
No package 'rdkafka' found
pkg-config: exit status 1
Run Code Online (Sandbox Code Playgroud)

我该如何解决 ?我尝试将它添加到路径中,但没有骰子!有什么建议 ?

go apache-kafka confluent-platform

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

Confluence 平台还是 Solace?

我们的企业同时具备 Solace 和 Confluence Platform 能力。

虽然 Solace 还支持实时流媒体和基于设备的产品,但企业为什么以及何时应该使用 Confluence 平台?

apache-kafka solace confluent-platform

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

Kafka Rest API KSQL查询永远等待并挂起

我正在使用从 curl POST 发布到 ksql 的 Kafka REST API 如果我不使用 LIMIT20,它会挂起。此外,如果我再次使用它来查询表,它会挂起。我从 python 脚本内部运行它在这里我在行时间之间查询 bcoz 我无法从流中获取最新结果,因为它是连续和持久的。

data = {"ksql":"SELECT MAX(ROWTIME),TIMESTAMPTOSTRING(ROWTIME, 'yyyy-MM-dd HH:mm:ss'),MYFIRMWAREVERSION,MYBASEMACID,BOOTTS,IMEI,PRODDEVICESERIALNUM,RESETREASON FROM NOV_STREAM WHERE TIMESTAMPTOSTRING(ROWTIME, 'yyyy-MM-dd HH:mm:ss') >= '2018-12-11 00:29:30'AND TIMESTAMPTOSTRING(ROWTIME, 'yyyy-MM-dd HH:mm:ss') <= '2018-12-11 23:29:30' GROUP BY ROWTIME,MYFIRMWAREVERSION,MYBASEMACID,BOOTTS,IMEI,PRODDEVICESERIALNUM,RESETREASON LIMIT 20;","streamsProperties":{"ksql.streams.auto.offset.reset": "earliest","format": "json"}}
Run Code Online (Sandbox Code Playgroud)

rest apache-kafka confluent-platform ksqldb

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

Confluence-Kafka Python:如何以编程方式列出所有主题

我的本地系统中运行着一个卡夫卡代理。为了使用我的基于 Django 的 Web 应用程序与损坏的设备进行通信,我使用confluence-kafka包装器。然而,通过浏览 admin api,我找不到任何用于列出 kafka 主题的 api。(主题是务实创建的并且是动态的)。

有什么办法可以在我的程序中列出它们吗?要求是,如果我的工作人员重新启动所有分配的监听这些主题的消费者,则必须重新初始化,因此我想循环到所有主题并为每个主题分配一个消费者。

python django python-3.x apache-kafka confluent-platform

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

类型错误:“mappingproxy”类型的对象不是 JSON 可序列化的

我尝试使用模式注册使用 confluent-kafka-python 的 AvroProducer 发布一条 avro 消息。但是代码无法序列化枚举类型。下面是代码和错误跟踪。任何帮助深表感谢。

from confluent_kafka import avro
from confluent_kafka.avro import AvroProducer
from example_schema.schema_classes import SCHEMA as value_schema
from example_schema.com.acme import *
import json

def function():
    avroProducer = AvroProducer({ 'bootstrap.servers': 'localhost:9092', 'schema.registry.url': 'http://localhost:8081' },  default_value_schema=value_schema)
    print(avroProducer)
    obj = Test()
    obj.name = 'vinay'
    obj.age = 11
    obj.sex = 'm'
    obj.myenum = Suit.CLUBS
    print(str(obj))
    avroProducer.produce(topic='test_topic',value=obj)
    avroProducer.flush()

function()

  File "main.py", line 16, in function
    avroProducer.produce(topic='test_topic',value=json.dumps(obj))
  File "/home/priv/anaconda3/lib/python3.6/site-packages/confluent_kafka/avro/__init__.py", line 80, in produce
    value = self._serializer.encode_record_with_schema(topic, value_schema, value)
  File "/home/priv/anaconda3/lib/python3.6/site-packages/confluent_kafka/avro/serializer/message_serializer.py", line 105, …
Run Code Online (Sandbox Code Playgroud)

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

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

无法反序列化主题数据

我正在使用这里的汇合 cp-all-in-one 项目配置:https://github.com/confluenceinc/cp-docker-images/blob/5.2.2-post/examples/cp-all-in-one /docker-compose.yml

http://localhost:8082/topics/zuum-positions 我正在使用以下 AVRO 正文发布一条消息:

{  
   "key_schema": "{\"type\":\"string\"}",
   "value_schema":"{  \"type\":\"record\",\"name\":\"Position\",\"fields\":[  {  \"name\":\"loadId\",\"type\":\"double\"},{\"name\":\"lat\",\"type\":\"double\"},{  \"name\":\"lon\",\"type\":\"double\"}]}",
   "records":[  
      {  
         "key":"22",
         "value":{  
            "lat":43.33,
            "lon":43.33,
            "loadId":22
         }
      }
   ]
}
Run Code Online (Sandbox Code Playgroud)

我已将以下标头正确添加到上述 POST 请求中: Content-Type: application/vnd.kafka.avro.v2+json Accept: application/vnd.kafka.v2+json

执行此请求时,我在 docker 日志中看到以下异常:

Error encountered in task zuum-sink-positions-0. Executing stage 'VALUE_CONVERTER' with class 'io.confluent.connect.avro.AvroConverter', where consumed record is {topic='zuum-positions', partition=0, offset=25, timestamp=1563480487456, timestampType=CreateTime}. org.apache.kafka.connect.errors.DataException: Failed to deserialize data for topic zuum-positions to Avro: 
connect            |    at io.confluent.connect.avro.AvroConverter.toConnectData(AvroConverter.java:107)
connect            |    at org.apache.kafka.connect.runtime.WorkerSinkTask.lambda$convertAndTransformRecord$1(WorkerSinkTask.java:487)
connect            |    at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:128) …
Run Code Online (Sandbox Code Playgroud)

apache-kafka docker confluent-schema-registry kafka-rest confluent-platform

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