- html - 出于某种原因,IE8 对我的 Sass 文件中继承的 html5 CSS 不友好?
- JMeter 在响应断言中使用 span 标签的问题
- html - 在 :hover and :active? 上具有不同效果的 CSS 动画
- html - 相对于居中的 html 内容固定的 CSS 重复背景?
我在 python 中一起玩 spark-streaming 和 kafka,并松散地跟随 this post但我对前面提到的 KafkaUtils.createStream() 函数有点困惑。
documentation通过明确解释主题词典的影响并没有做太多事情。但我怀疑我之所以这么认为,是因为我对 kafka 的工作原理了解不深,而答案是显而易见的。
我知道它应该是这样的字典:{"topic.name": 1}
我可以重复文档并说这意味着创建的流将从单个分区中消耗。
所以我想我只是在寻找关于这个特定功能的用法以及我对 kafka 概念的理解的一些说明。我们将使用以下示例:
假设我已经定义了一个主题 my.topic
,它有 3 个分区,其传入消息按一个键拆分,假设是一个用户 ID。
如果我像这样初始化一个流:
from pyspark.streaming.kafka import KafkaUtils
kafkaStream = KafkaUtils.createStream(
ssc,
'kafka:2181',
'consumer-group-name',
{'my.topic':1}
)
我认为这个流只会从一个分区消费,所以不会看到进入 my.topic
的每条消息,我的想法是否正确?换句话说,它只会看到从 userid 发送到 3 个分区之一的消息?
我的问题是:
如何正确设置此参数以使用发送到 my.topic
的所有消息?
我的直觉是我只需将主题参数设置为 {'my.topic': 3}
,那么我的问题就变成了:
为什么我会使用小于分区总数的数字?
我的直觉告诉我,这与您所做的工作的“原子性”程度有关。例如,如果我只是简单地转换数据(比如,从 CSV 到 JSON 文档列表或其他东西)然后将上面的 3 个流都设置为 {'my.topic': 1}
它们的主题参数和同一消费者组的所有部分将有利于从每个分区启用并行消费,因为不需要共享有关消费的每条消息的信息。
与此同时,如果我计算的是整个主题 I.E. 的实时指标。带过滤器的时间窗口平均值等。我很难找到一种方法来实现类似的东西而不设置消费者组 I.E. 中每个分量信号的更复杂的下游处理Sum1 + Sum2 + Sum3 = 总和
但我的知识还是处于使用 Kafka 和 Spark 的“初级”阶段。
有没有办法告诉 createStream() 使用所有分区,而无需提前知道有多少分区?类似于 {'my.topic': -1}
?
是否可以在一个流中指定多个主题? IE。 {'my.topic': 1, 'my.other.topic': 1}
我真的很讨厌这个问题的答案只是“是的,你的直觉是正确的。”。最好的情况是有人告诉我我误解了所有事情并让我直截了当。所以请...这样做吧!
最佳答案
这是 Kafka-Spark 集成页面中提到的内容。
val kafkaStream = KafkaUtils.createStream(streamingContext, [ZK quorum], [consumer group id], [per-topic number of Kafka partitions to consume])
KafkaUtils.createStream 将创建一个接收器并使用 Kafka 主题。
“要使用的每个主题的 Kafka 分区数”选项表示此接收器将并行读取多少个分区。
例如,您有一个名为“Topic1”的主题,有 2 个分区,并且您提供了选项“Topic1”:1,那么 Kafka 接收器将一次读取 1 个分区 [它最终会读取所有分区,但会读取一次一个分区]。这样做的原因是读取分区中的消息并保留数据写入主题的顺序。
例如,Topic1 的分区 1 包含消息 {1,11,21,31,41},分区 2 包含消息 {2,12,22,32,42},然后使用上述设置读取将产生类似 { 1,11,21,31,41,2,12,22,32,42}。每个分区中的消息是单独读取的,因此不会混合在一起。
如果您提供的选项为“Topic1”:2,那么接收方将一次读取 2 个分区,并且这些分区中的消息将混合在一起。对于上面相同的启动示例,具有“Topic1”的接收者:2 将产生类似于 {1,2,11,12,21,22....}
将此视为接收器可以对给定主题分区执行的并行读取数。
<强>5。一个流中可以指定多个主题吗?是的你可以。
关于python - 在 KafkaUtils.createstream() 中使用 "topics"参数的正确方法是什么?,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/48161253/
简而言之:我想从可变参数模板参数中提取各种选项,但不仅通过标签而且通过那些参数的索引,这些参数是未知的 标签。我喜欢 boost 中的方法(例如 heap 或 lockfree 策略),但想让它与 S
我可以对单元格中的 excel IF 语句提供一些帮助吗? 它在做什么? 对“BaselineAmount”进行了哪些评估? =IF(BaselineAmount, (Variance/Baselin
我正在使用以下方法: public async Task Save(Foo foo,out int param) { ....... MySqlParameter prmparamID
我正在使用 CodeGear RAD Studio IDE。 为了使用命令行参数测试我的应用程序,我多次使用了“运行 -> 参数”菜单中的“参数”字段。 但是每次我给它提供一个新值时,它都无法从“下拉
我已经为信用卡类编写了一些代码,粘贴在下面。我有一个接受上述变量的构造函数,并且正在研究一些方法将这些变量格式化为字符串,以便最终输出将类似于 号码:1234 5678 9012 3456 截止日期:
MySql IN 参数 - 在存储过程中使用时,VarChar IN 参数 val 是否需要单引号? 我已经像平常一样创建了经典 ASP 代码,但我没有更新该列。 我需要引用 VarChar 参数吗?
给出了下面的开始,但似乎不知道如何完成它。本质上,如果我调用 myTest([one, Two, Three], 2); 它应该返回元素 third。必须使用for循环来找到我的解决方案。 funct
将 1113355579999 作为参数传递时,该值在函数内部变为 959050335。 调用(main.c): printf("%d\n", FindCommonDigit(111335557999
这个问题在这里已经有了答案: Is Java "pass-by-reference" or "pass-by-value"? (92 个回答) 关闭9年前。 public class StackOve
我真的很困惑,当像 1 == scanf("%lg", &entry) 交换为 scanf("%lg", &entry) == 1 没有区别。我的实验书上说的是前者,而我觉得后者是可以理解的。 1 =
我正在尝试使用调用 SetupDiGetDeviceRegistryProperty 的函数使用德尔福 7。该调用来自示例函数 SetupEnumAvailableComPorts .它看起来像这样:
我需要在现有项目上实现一些事件的显示。我无法更改数据库结构。 在我的 Controller 中,我(从 ajax 请求)传递了一个时间戳,并且我需要显示之前的 8 个事件。因此,如果时间戳是(转换后)
rails 新手。按照多态关联的教程,我遇到了这个以在create 和destroy 中设置@client。 @client = Client.find(params[:client_id] || p
通过将 VM 参数设置为 -Xmx1024m,我能够通过 Eclipse 运行 Java 程序-Xms256M。现在我想通过 Windows 中的 .bat 文件运行相同的 Java 程序 (jar)
我有一个 Delphi DLL,它在被 Delphi 应用程序调用时工作并导出声明为的方法: Procedure ProduceOutput(request,inputs:widestring; va
浏览完文档和示例后,我还没有弄清楚 schema.yaml 文件中的参数到底用在哪里。 在此处使用 AWS 代码示例:https://github.com/aws-samples/aws-proton
程序参数: procedure get_user_profile ( i_attuid in ras_user.attuid%type, i_data_group in data_g
我有一个字符串作为参数传递给我的存储过程。 dim AgentString as String = " 'test1', 'test2', 'test3' " 我想在 IN 中使用该参数声明。 AND
这个问题已经有答案了: When should I use "this" in a class? (17 个回答) 已关闭 6 年前。 我运行了一些java代码,我看到了一些我不太明白的东西。为什么下
我输入 scroll(0,10,200,10);但是当它运行时,它会传递字符串“xxpos”或“yypos”,我确实在没有撇号的情况下尝试过,但它就是行不通。 scroll = function(xp
我是一名优秀的程序员,十分优秀!