当我在没有 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
我正在尝试使用 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) 在机器A 中kafka 正在运行,在机器B中安装了 hadoop。现在我想从 kafka 向 hadoop 写入数据。我已经在 Machine A 中安装了 Confluent Platform 。
任何人都可以指导我必须添加哪些配置才能将数据从 kafka 写入在不同机器上运行的 hadoop
我正在用 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
在 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)
我该如何解决 ?我尝试将它添加到路径中,但没有骰子!有什么建议 ?
我们的企业同时具备 Solace 和 Confluence Platform 能力。
虽然 Solace 还支持实时流媒体和基于设备的产品,但企业为什么以及何时应该使用 Confluence 平台?
我正在使用从 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) 我的本地系统中运行着一个卡夫卡代理。为了使用我的基于 Django 的 Web 应用程序与损坏的设备进行通信,我使用confluence-kafka包装器。然而,通过浏览 admin api,我找不到任何用于列出 kafka 主题的 api。(主题是务实创建的并且是动态的)。
有什么办法可以在我的程序中列出它们吗?要求是,如果我的工作人员重新启动所有分配的监听这些主题的消费者,则必须重新初始化,因此我想循环到所有主题并为每个主题分配一个消费者。
我尝试使用模式注册使用 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
我正在使用这里的汇合 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
apache-kafka ×9
avro ×3
docker ×2
python ×2
django ×1
go ×1
hadoop ×1
hdfs ×1
java ×1
kafka-rest ×1
ksqldb ×1
python-3.x ×1
rest ×1
solace ×1