我想让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上的防火墙?
谢谢
我想使用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) 我正在通过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启动它.
有些人可以通过一些亮点甚至指向一些资源也会有所帮助.
谢谢.
我想创建一个并发 @KafkaListener,它可以处理多个主题,每个主题都有不同数量的分区。
我注意到 Spring-Kafka 只为大多数分区的主题为每个分区初始化一个使用者。
示例:我将并发设置为 8。我@KafkaListener听了以下主题。主题 A 的分区最多 - 5 个,因此 Spring-Kafka 初始化了 5 个消费者。我希望 Spring-Kafka 初始化 8 个消费者,这是根据我的并发属性允许的最大值。
不初始化更多消费者的技术原因是什么?
我如何绕过这个,以便我可以使用@KafkaListener注释初始化更多的使用者?(如果可能的话)
我想用 python 获取 Kafka 消费者组列表,但我不能。
我使用zookeeper python客户端(kazoo),但消费者组列表为空,因为这种方法适用于旧消费者,而我们没有使用旧消费者。
如何使用python代码获取消费者组列表?
./kafka-consumer-groups.sh -bootstrap-server localhost:9092 -list
Run Code Online (Sandbox Code Playgroud) 先将数据放入 Kafka,然后再放入数据库,或者其他方式,优点和缺点是什么?
示例:用户执行 REST (POST) 调用来存储产品。通常我会在后端接听这个调用并将正文保存到数据库中(在验证之后......)。最佳实践是接听此调用并将数据存储在 Kafka 中,然后将其保存到数据库(在本例中,数据库是 kafka 消费者)。
还是先保存到数据库,然后发送到kafka比较好?
谢谢
我的系统上安装了 Acanonda。我通过提供安装了tensorflow
pip install tensorflow
Run Code Online (Sandbox Code Playgroud)
它已成功安装:
以下是最后的跟踪:
Successfully installed absl-py-0.7.0 astor-0.7.1 gast-0.2.2 grpcio-1.19.0
keras-applications-1.0.7 keras-preprocessing-1.0.9 markdown-3.0.1 mock-2.0.0
pbr-5.1.3 protobuf-3.7.0 tensorboard-1.13.0 tensorflow-1.13.1 tensorflow-
estimator-1.13.0 termcolor-1.1.0
Run Code Online (Sandbox Code Playgroud)
现在我尝试运行以下脚本。没什么特别的,只是导入库。
import tensorflow as tf
from tensorflow import keras
import numpy as np
import matplotlib.pyplot as plt
Run Code Online (Sandbox Code Playgroud)
在运行脚本时,我收到以下错误:
ModuleNotFoundError: No module named 'numpy.core._multiarray_umath'
ImportError: numpy.core.multiarray failed to import
The above exception was the direct cause of the following exception:
Traceback (most recent call last):
File "<frozen importlib._bootstrap>", line 968, in _find_and_load
SystemError: <class '_frozen_importlib._ModuleLockManager'> returned a …Run Code Online (Sandbox Code Playgroud) 我正在尝试配置 kafka 客户端以针对安全的 kafka 服务器进行身份验证。我已经设置了 jaas 和 ssl 配置,但它在抱怨 serviceNames。
我没有使用 Kerberos。
命令
KAFKA_OPTS="-Djava.security.auth.login.config=./jaas.conf" \
kafka-console-producer --broker-list k0:9092,k1:9092,k2:9092 \
--topic test-topic
--producer.config ./ssl.properties
Run Code Online (Sandbox Code Playgroud)
错误
org.apache.kafka.common.KafkaException: Failed to construct kafka producer
at org.apache.kafka.clients.producer.KafkaProducer.<init>
[ ... ]
Caused by: java.lang.IllegalArgumentException: No serviceName defined in either JAAS or Kafka config
Run Code Online (Sandbox Code Playgroud)
配置文件
KafkaServer {
org.apache.kafka.common.security.plain.PlainLoginModule required
serviceName="kafka"
password="broker-secret"
user_broker="broker-secret"
sasl.enabled.mechanisms=PLAIN
sasl.mechanism.inter.broker.protocol=PLAIN
confluent.metrics.reporter.sasl.mechanism=PLAIN
user_username1="password1";
};
Run Code Online (Sandbox Code Playgroud)
ssl.properties
bootstrap.servers=k0:9092,k1:9092,k2:9092
security.protocol=SASL_PLAINTEXT
ssl.truststore.location=/var/ssl/private/client.truststore.jks
ssl.truststore.password=confluent
ssl.keystore.location=/var/ssl/private/client.keystore.jks
ssl.keystore.password=confluent
ssl.key.password=confluent
producer.bootstrap.servers=k0:9092,1:9092,k2:9092
producer.security.protocol=SASL_PLAINTEXT
producer.ssl.truststore.location=/var/private/ssl/kafka.client.truststore.jks
producer.ssl.truststore.location=/var/ssl/private/client.truststore.jks
producer.ssl.truststore.password=confluent
producer.ssl.keystore.location=/var/ssl/private/client.keystore.jks
producer.ssl.keystore.password=confluent
producer.ssl.key.password=confluent
org.apache.kafka.common.security.plain.PlainLoginModule required
password="broker-secret"
user_broker="broker-secret" …Run Code Online (Sandbox Code Playgroud) 有没有办法在 Kafka Connect 启动时自动加载(多个)Kafka Connect 连接器(例如在 Confluence Platform 中)?
到目前为止我发现了什么:
Confluence Docs 声明使用bin/connect-standalone
独立模式的命令以及工作线程和每个连接器的属性文件。
对于分布式模式,您必须通过 REST API 运行连接器。
https://docs.confluence.io/current/connect/userguide.html#standalone-mode,https://docs.confluence.io/current/connect/managing/configuring.html#standalone-example
是否有另一种方法,例如包含应在“connect-[standalone|distributed].properties”文件中运行的所有连接器(类似于在 ksql-server.properties 中提供 KSQL 查询文件),以便它们自动加载到Kafka Connect 的启动(例如在 Confluence 平台中)?
或者,即使在生产环境中,连接器是否也如上所述“手动”加载?
我正在使用 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
apache-kafka ×9
python ×3
confluent ×2
database ×1
jdbc ×1
kafka-python ×1
kafka-topic ×1
mysql ×1
numpy ×1
security ×1
spring-kafka ×1
ssl ×1
tensorflow ×1