我使用控制台脚本与使用sasl保护的Kafka进行通信时遇到问题.Kafka用sasl保护,监听器是SASL_PLAINTEXT,机制是PLAIN.
我做了什么:我尝试使用一个kafka脚本列出一些数据:
bin/kafka-consumer-groups.sh --bootstrap-server(address)--list
但是我得到" WARN Bootstrap broker(地址)已断开连接(org.apache.kafka.clients.NetworkClient) "并且命令失败,这是可以理解的,因为它是用sasl保护的.
所以我尝试了如何在该命令中添加客户端用户名/密码.首先,我尝试运行kafka-console-consumer脚本,我使用--command-config添加必要的文件.我很快发现我不能直接添加jaas文件,我需要使用.properties文件,所以我做到了.
我的属性文件(请记住括号表示"删失"数据,我不能在这里放置所有真实数据):
bin/kafka-consumer-groups.sh --bootstrap-server (address) --list
Run Code Online (Sandbox Code Playgroud)
我的jaas文件:
WARN Bootstrap broker (address) disconnected (org.apache.kafka.clients.NetworkClient)
Run Code Online (Sandbox Code Playgroud)
这个jaas文件适用于我的标准java应用程序.
但是,当我尝试运行kafka-consumer-groups脚本或kafka-console-consumer时,我收到此错误:
bootstrap.servers=(address)
zookeeper.connect=127.0.0.1:2181
zookeeper.connection.timeout.ms=6000
sasl.jaas.config=(path)/consumer_jaas.conf
security.protocol=SASL_PLAINTEXT
sasl.mechanism=PLAIN
group.id=(group)
Run Code Online (Sandbox Code Playgroud)
这个jaas文件是我在java应用程序中使用的文件的直接副本,它与kafka通信并且它可以工作,但是在这里,使用控制台工具,它只是不起作用.我试着寻找解决方案,但我找不到任何有用的东西.
谁能帮我这个?
我想用种子np.random.permutation,比如
np.random.permutation(10, seed=42)
Run Code Online (Sandbox Code Playgroud)
我收到以下错误:
"permutation() takes no keyword arguments"
Run Code Online (Sandbox Code Playgroud)
我该怎么办呢?谢谢.
Kafka的主要开发人员表示Kafka是CA而CAP是CAP定理.但我很困惑,卡夫卡不是分区容忍的吗?我认为确实如此,当一个复制失败时,另一个将成为领导者并继续工作!
另外,我想知道如果Kafka使用P怎么办?P会伤害C还是A?
我是hadoop的新手.我正在尝试在hdfs中创建一个目录,但我无法创建.
我已登录"hduser"因此我认为/ home/hduser"预先存在为Unix fs.所以我尝试使用下面的命令创建hadoop目录.
[hduser@Virus ~]$ hadoop fs -mkdir /home/hduser/mydata/
14/12/03 15:04:53 WARN util.NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
mkdir: `/home/hduser/mydata/': No such file or directory
Run Code Online (Sandbox Code Playgroud)
在线搜索后,我想到hadoop可能无法理解"/ home/hduser"或者我使用hadoop2,其中mkdir不会像Unix命令"madir -p"那样工作(递归).因此我试图创建"/ mydata"但没有运气.
[hduser@Virus ~]$ hadoop fs -mkdir /mydata
14/12/03 15:09:26 WARN util.NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
mkdir: Cannot create directory /mydata. Name node is in safe mode.
Run Code Online (Sandbox Code Playgroud)
我试图离开安全模式,但仍然存在问题.
[hduser@Virus ~]$ hdfs dfsadmin -safemode leave
14/12/03 …Run Code Online (Sandbox Code Playgroud) 我试图了解如何使用微服务(使用 apache kafka)在事件驱动的世界中处理异常。例如,如果您采用以下订单场景,其中需要在完成订单之前执行以下操作。
在这种情况下的任何阶段,都可能出现故障,例如:
您如何跟踪每个阶段已被要求和/或完成?
你如何处理出现的问题?你将如何通知前端失败?
我们遇到了卡夫卡的问题.有时突然间,我们会在没有警告的情况下退出同步并在发出事件时开始获取异常.
我们得到的例外是:"java.io.IOException:打开的文件过多"
在许多情况下,这似乎是kafka抛出的一般异常.我们稍微调查一下,我们认为根本原因是当尝试向某个主题发出事件时,它会失败,因为kafka没有针对此主题的领导分区
有人可以帮忙吗?
我有一个数据流用例,我希望根据每个客户存储库(可能大约100,000个)定义主题.每个数据流都是一个带有分区的主题(大约几十个)定义流程的不同阶段.
卡夫卡是否适合这样的场景?如果不是,我将如何改造我的用例来处理这种情况.此外,即使在处理过程中,每个客户存储库数据也不能与其他客户存储库数据混合.
data.rdd.getNumPartitions() # output 2456
Run Code Online (Sandbox Code Playgroud)
然后我这样做
data.rdd.repartition(3000)
但
data.rdd.getNumPartitions()#outout仍然是2456
如何更改分区数量.一种方法可以是首先将DF转换为rdd,重新分区然后将rdd转换回DF.但这需要很多时间.越来越多的分区是否使操作更加分散,因此更快?谢谢
machine-learning bigdata apache-spark apache-spark-sql pyspark
我最近启用了 SSL 并尝试以分布式模式启动 Kafka 连接。跑步时
connect-distributed connect-distributed.properties
Run Code Online (Sandbox Code Playgroud)
我收到以下错误:
[2018-10-09 16:50:57,190] INFO Stopping task (io.confluent.connect.jdbc.sink.JdbcSinkTask:106)
[2018-10-09 16:50:55,471] ERROR WorkerSinkTask{id=sink-mariadb-test} Task threw an uncaught and unrecoverable exception (org.apache.kafka.connect.runtime.WorkerTask:177)
java.lang.OutOfMemoryError: Java heap space
at java.nio.HeapByteBuffer.<init>(HeapByteBuffer.java:57)
at java.nio.ByteBuffer.allocate(ByteBuffer.java:335)
at org.apache.kafka.common.memory.MemoryPool$1.tryAllocate(MemoryPool.java:30)
at org.apache.kafka.common.network.NetworkReceive.readFrom(NetworkReceive.java:112)
at org.apache.kafka.common.network.KafkaChannel.receive(KafkaChannel.java:344)
at org.apache.kafka.common.network.KafkaChannel.read(KafkaChannel.java:305)
at org.apache.kafka.common.network.Selector.attemptRead(Selector.java:560)
at org.apache.kafka.common.network.Selector.pollSelectionKeys(Selector.java:496)
at org.apache.kafka.common.network.Selector.poll(Selector.java:425)
at org.apache.kafka.clients.NetworkClient.poll(NetworkClient.java:510)
at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:271)
at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:242)
at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:218)
at org.apache.kafka.clients.consumer.internals.AbstractCoordinator.ensureCoordinatorReady(AbstractCoordinator.java:230)
at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.poll(ConsumerCoordinator.java:314)
at org.apache.kafka.clients.consumer.KafkaConsumer.updateAssignmentMetadataIfNeeded(KafkaConsumer.java:1218)
at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1181)
at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1115)
at org.apache.kafka.connect.runtime.WorkerSinkTask.pollConsumer(WorkerSinkTask.java:444)
at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:317)
at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:225)
at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:193)
at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:175)
at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:219)
at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) …Run Code Online (Sandbox Code Playgroud) 这是我在 Stackoverflow 上的第一篇文章,希望我没有选错部分。
语境 :
Kafka HEAP 大小在以下文件中配置:
/etc/systemd/system/kafka.service
Run Code Online (Sandbox Code Playgroud)
使用以下参数:
Environment="KAFKA_HEAP_OPTS=-Xms6g -Xmx6g"
Run Code Online (Sandbox Code Playgroud)
操作系统是“CentOS Linux 7.7.1908 版”。
Kafka 是"confluent-kafka-2.12-5.3.1-1.noarch",从以下存储库安装:
# Confluent REPO
[Confluent.dist]
name=Confluent repository (dist)
baseurl=http://packages.confluent.io/rpm/5.3/7
gpgcheck=1
gpgkey=http://packages.confluent.io/rpm/5.3/archive.key
enabled=1
[Confluent]
name=Confluent repository
baseurl=http://packages.confluent.io/rpm/5.3
gpgcheck=1
gpgkey=http://packages.confluent.io/rpm/5.3/archive.key
enabled=1
Run Code Online (Sandbox Code Playgroud)
几天前,我在 3 台机器 KAFKA 集群上激活了 SSL,突然,以下命令停止工作:
kafka-topics --bootstrap-server <the.fqdn.of.server>:9093 --describe --topic <TOPIC-NAME>
Run Code Online (Sandbox Code Playgroud)
其中返回以下错误:
[2019-10-03 11:38:52,790] ERROR Uncaught exception in thread 'kafka-admin-client-thread | adminclient-1':(org.apache.kafka.common.utils.KafkaThread)
java.lang.OutOfMemoryError: Java heap space
at java.nio.HeapByteBuffer.<init>(HeapByteBuffer.java:57)
at java.nio.ByteBuffer.allocate(ByteBuffer.java:335)
at org.apache.kafka.common.memory.MemoryPool$1.tryAllocate(MemoryPool.java:30)
at org.apache.kafka.common.network.NetworkReceive.readFrom(NetworkReceive.java:112)
at org.apache.kafka.common.network.KafkaChannel.receive(KafkaChannel.java:424)
at org.apache.kafka.common.network.KafkaChannel.read(KafkaChannel.java:385)
at …Run Code Online (Sandbox Code Playgroud) apache-kafka ×6
java ×2
apache-spark ×1
bigdata ×1
event-driven ×1
hadoop ×1
hdfs ×1
heap-memory ×1
jaas ×1
numpy ×1
permutation ×1
pyspark ×1
python ×1
random ×1
sasl ×1
shell ×1
ssl ×1