小编Gio*_*ous的帖子

Kafka"未在JAAS配置中指定登录模块"

我使用控制台脚本与使用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通信并且它可以工作,但是在这里,使用控制台工具,它只是不起作用.我试着寻找解决方案,但我找不到任何有用的东西.

谁能帮我这个?

java sasl jaas apache-kafka

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

np.random.permutation与种子?

我想用种子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)

我该怎么办呢?谢谢.

python random numpy permutation

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

为什么卡夫卡不是CAP定理中的P.

Kafka的主要开发人员表示Kafka是CA而CAP是CAP定理.但我很困惑,卡夫卡不是分区容忍的吗?我认为确实如此,当一个复制失败时,另一个将成为领导者并继续工作!

另外,我想知道如果Kafka使用P怎么办?P会伤害C还是A?

apache-kafka

11
推荐指数
3
解决办法
2911
查看次数

在hadoop文件系统中创建目录

我是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)

shell hadoop command-line-interface hdfs

10
推荐指数
3
解决办法
3万
查看次数

在事件驱动的世界中处理异常

我试图了解如何使用微服务(使用 apache kafka)在事件驱动的世界中处理异常。例如,如果您采用以下订单场景,其中需要在完成订单之前执行以下操作。

  • 1) 向支付服务提供商授权支付
  • 2)从库存中保留该项目
  • 3.1) 通过支付服务提供商捕获付款
  • 3.2) 订购商品
  • 4) 发送电子邮件通知接受订单并附上收据

在这种情况下的任何阶段,都可能出现故障,例如:

  • 该商品不再有库存
  • 支付信息有误
  • 收款人使用的账户没有可用资金
  • 外部调用(例如对支付服务提供商的调用)失败,例如停机

您如何跟踪每个阶段已被要求和/或完成?

你如何处理出现的问题?你将如何通知前端失败?

event-driven-design event-driven

10
推荐指数
1
解决办法
1944
查看次数

kafka同步:"java.io.IOException:太多打开的文件"

我们遇到了卡夫卡的问题.有时突然间,我们会在没有警告的情况下退出同步并在发出事件时开始获取异常.

我们得到的例外是:"java.io.IOException:打开的文件过多"

在许多情况下,这似乎是kafka抛出的一般异常.我们稍微调查一下,我们认为根本原因是当尝试向某个主题发出事件时,它会失败,因为kafka没有针对此主题的领导分区

有人可以帮忙吗?

apache-kafka kafka-consumer-api

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

我可以在Kafka Cluser中拥有数以千计的主题吗?

我有一个数据流用例,我希望根据每个客户存储库(可能大约100,000个)定义主题.每个数据流都是一个带有分区的主题(大约几十个)定义流程的不同阶段.

卡夫卡是否适合这样的场景?如果不是,我将如何改造我的用例来处理这种情况.此外,即使在处理过程中,每个客户存储库数据也不能与其他客户存储库数据混合.

apache-kafka

9
推荐指数
1
解决办法
6360
查看次数

如何重新分区pyspark数据帧?

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

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

启用 SSL 后 Kafka Connect 耗尽 Java 堆空间

我最近启用了 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)

java heap-memory apache-kafka apache-kafka-connect

9
推荐指数
1
解决办法
6635
查看次数

KAFKA 和 SSL:java.lang.OutOfMemoryError:在 KAFKA SSL 集群上使用 kafka-topics 命令时的 Java 堆空间

这是我在 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)

ssl apache-kafka apache-kafka-security

9
推荐指数
1
解决办法
5421
查看次数