小编Gio*_*ous的帖子

在单个节点上将SSL与Kafka一起使用

我想让SSL与Kafka一起运行,以使其更加安全。我下载并安装了Kafka。我完全按照说明为SSL创建证书和信任库。我将以下内容添加到我的config / server.properties中

ssl.enabled.protocols=TLSv1.2,TLSv1.1,TLSv1
ssl.keystore.type=JKS
ssl.truststore.type=JKS

listeners=PLAINTEXT://localhost:9092,SSL://localhost:9093
ssl.endpoint.identification.algorithm=HTTPS
security.inter.broker.protocol=SSL
ssl.secure.random.implementation=SHA1PRNG
ssl.endpoint.identification.algorithm=HTTPS
ssl.keystore.location=/home/ec2-user/workspace/kafka/cert/server.keystore.jks
ssl.key.password=<the password>
ssl.keystore.password=<the password>
ssl.truststore.location=/home/ec2-user/workspace/kafk/cert/server.truststore.jks
ssl.truststore.password=<the password>
Run Code Online (Sandbox Code Playgroud)

启动Zookeeper之后,启动kafak时出现此错误:[2017-12-07 16:02:52,155]错误[Controller id = 0,targetBrokerId = 0]由于以下原因,与节点0的连接身份验证失败:SSL握手失败( org.apache.kafka.clients.NetworkClient)。我必须终止任务才能停止此消息

看logs/controller.log:

[Controller-0-to-broker-0-send-thread]: Controller 0's connection to broker localhost:9093 (id: 0 rack: null) was unsuccessful (kafka.controller.RequestSendThread)
Run Code Online (Sandbox Code Playgroud)

您是否必须打开端口9093上的防火墙?

谢谢

security ssl apache-kafka

5
推荐指数
1
解决办法
2863
查看次数

生产中的Kafka监控工具

需要检查生产中监控Kafka的工具.此外,工具不需要许可证或重型硬件.特别是我需要一个工具来评估消费者对主题的偏差,主题的健康状况.

monitoring apache-kafka

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

生产者如何找到 kafka 阅读器

生产者通过设置 Kafka Broker 列表来发送消息,如下所示。

props.put("bootstrap.servers", "127.0.0.1:9092,127.0.0.1:9092,127.0.0.1:9092");
Run Code Online (Sandbox Code Playgroud)

我想知道“生产者”如何知道三个经纪人中的哪一个知道哪个有分区领导者。对于典型的分布式服务器,要么你有一个承载服务器,要么有一个虚拟IP,但对于Kafka,它是如何加载的?生产者程序是否尝试随机连接到一个代理并寻找具有分区领导者的代理?

apache-kafka

5
推荐指数
1
解决办法
1426
查看次数

如何为新的Kafka消费者设置consumer.id

根据文档,Old Consumer Configsconsumer.id默认包含一个为 null的集合:

consumer.id
默认值:null
描述:如果未设置则自动生成。

是否可以consumer.id为New Kafka Consumer 设置 ,如果可以,我该如何实现?

apache-kafka kafka-consumer-api

5
推荐指数
1
解决办法
1046
查看次数

数据在 kafka 服务器中存储多长时间?

假设我要从 kafka 生产者向 kafka 消费者发送一些消息,那么它将存储在哪里?是否有用于存储消息的数据库?消息存储多长时间?

任何人都可以请解释一下。

apache-kafka

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

无法使用Confluent REST API运行JDBC Source连接器

我想使用Confluent的REST API运行JDBC源连接器.虽然独立模式使用以下属性文件可以完美运行:

name=source-mysql-test
connector.class=io.confluent.connect.jdbc.JdbcSourceConnector
tasks.max=1

connection.url=jdbc:mysql://localhost:3306/kafka
connection.user=myuser
connection.password=mypass


table.whitelist=MY_TABLE

# Pull all rows based on timestamp
mode=timestamp
timestamp.column.name=ROWVERSION
validate.non.null=false

# The Kafka topic will be made up of this prefix, plus the table name.
topic.prefix=MYSQL-

table.types=TABLE,VIEW
poll.interval.ms=1000
Run Code Online (Sandbox Code Playgroud)

我无法使用REST API运行连接器.这是电话:

curl -X POST -H "Content-Type: application/json" --data '{"name": "source-mysql-test", "config": {"connector.class":"io.confluent.connect.jdbc.JdbcSourceConnector", "tasks.max":"1", "connection.url":"jdbc:mysql://localhost:3306/kafka","connection.user":"myuser","connection.password":"mypass", "table.whitelist":"MY_TABLE", "mode":"timestamp", "timestamp.column.name":"ROWVERSION", "validate.non.null":"false", "topic.prefix":"MYSQL-", "table.types":"TABLE,VIEW", "poll.interval.ms":"1000" }}' http://localhost:8083/connectors
Run Code Online (Sandbox Code Playgroud)

以下是回复:

{
  "error_code": 400,
  "message": "Connector configuration is invalid and contains the following 2 error(s):\nInvalid value com.mysql.jdbc.exceptions.jdbc4.MySQLNonTransientConnectionException: Could not …
Run Code Online (Sandbox Code Playgroud)

mysql jdbc apache-kafka confluent apache-kafka-connect

5
推荐指数
1
解决办法
798
查看次数

Kafka连接群集设置或启动连接工作者

我正在通过kafka connect,我正在尝试获取这些概念.

让我们说我有kafka集群(节点k1,k2和k3)设置并且它正在运行,现在我想在不同的节点中运行kafka connect worker,比如分布式模式下的c1和c2.

几个问题.

1)要在分布式模式下运行或启动kafka connect,我需要使用../bin/connect-distributed.shkakfa集群节点中可用的命令,所以我需要从任何一个kafka集群节点启动kafka connect?或者我启动kafka connect的任何节点都需要有kafka二进制文件才能使用../bin/connect-distributed.sh

2)我需要将我的连接器插件复制到我执行第1步的任何kafka集群节点(或所有集群节点?)?

3)在工作节点上启动jvm进程之前,kafka如何将这些连接器插件复制到工作节点?因为插件是具有我的任务代码的插件,需要将其复制到worker才能在worker中启动进程.

4)我是否需要在连接群集节点c1和c2中安装任何东西,比如需要安装java或任何相关的kafka连接?

5)在某些地方,它说使用汇合平台,但我想先用apache kafka connect启动它.

有些人可以通过一些亮点甚至指向一些资源也会有所帮助.

谢谢.

apache-kafka confluent apache-kafka-connect

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

返回带有序列化程序列表的响应Django REST Framework

I'm coding some backend software for a second-hand selling app using Django and DjangoRestFramework. Right now, I'm trying to send a Response object that contains a list of products, but I seem not to be able to return an actual list of products, as I'm getting an error saying

ListSerializer is not JSON serializable.
Run Code Online (Sandbox Code Playgroud)

I've tried both using the serializer constructor like this:

ProductoSerializer(products, many=True)
Run Code Online (Sandbox Code Playgroud)

And by creating a list of ProductoSerializer.data and then creating the Response object with that. …

python django django-rest-framework

5
推荐指数
1
解决办法
243
查看次数

使用 Pandas query() 过滤时间戳列上的数据框

我正在尝试使用query()时间戳列上的字符串和函数过滤 Pandas 数据框:

df.query('Timestamp < "2020-02-01"')
Run Code Online (Sandbox Code Playgroud)

但是,我收到以下错误:

Traceback (most recent call last):   
File "C:\ENERCON\Python 3.7.2\lib\site-packages\IPython\core\interactiveshell.py", line 3326, in run_code
     exec(code_obj, self.user_global_ns, self.user_ns)   
File "<ipython-input-3-7bb40e9c631a>", line 1, in <module>
     df.query('Timestamp < "2020-02-01"')   
File "C:\ENERCON\Python 3.7.2\lib\site-packages\pandas\core\frame.py", line 3199, in query
     res = self.eval(expr, **kwargs)   
File "C:\ENERCON\Python 3.7.2\lib\site-packages\pandas\core\frame.py", line 3315, in eval
     return _eval(expr, inplace=inplace, **kwargs)   
File "C:\ENERCON\Python 3.7.2\lib\site-packages\pandas\core\computation\eval.py", line 327, in eval
     ret = eng_inst.evaluate()   
File "C:\ENERCON\Python 3.7.2\lib\site-packages\pandas\core\computation\engines.py", line 142, in evaluate
     return self.expr()   
File "C:\ENERCON\Python 3.7.2\lib\site-packages\pandas\core\computation\expr.py", line 837, in …
Run Code Online (Sandbox Code Playgroud)

python timestamp dataframe pandas

5
推荐指数
1
解决办法
1059
查看次数

Python-Kafka:无限轮询主题

我正在使用 python-kafka 来收听 kafka 主题并使用该记录。我想让它无限轮询而不退出。这是我的代码如下:

def test():
    consumer = KafkaConsumer('abc', 'localhost:9092', auto_offset_reset='earliest')
    for msg in consumer:
        print(msg.value)
Run Code Online (Sandbox Code Playgroud)

这段代码只是读取数据,直接退出。有没有办法即使没有推送消息也可以继续收听主题?

任何持续监控该主题的相关示例对我来说也很棒。

python apache-kafka kafka-consumer-api kafka-python kafka-topic

5
推荐指数
1
解决办法
8136
查看次数