标签: apache-samza

Apache Samza和Apache Storm在用例方面有何不同?

我偶然发现了这篇文章,声称将Samza与Storm进行了对比,但似乎只是为了解决实现细节问题.

这两个分布式计算引擎的用例在哪些方面有所不同?每种工具都适合做什么工作?

apache-storm apache-samza

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

Apache Storm vs Apache Samza vs Apache Spark

我曾经参与过Storm和Spark,但是Samza很新.

我不明白为什么当Storm已经在那里进行实时处理时Samza被引入了.Spark在内存中提供近实时处理,并具有其他非常有用的组件如graphx和mllib.

Samza带来了哪些改进以及可能的进一步改进?

apache-spark apache-storm apache-samza

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

Samza/Kafka无法更新元数据

我目前正在编写一个Samza脚本,它将从Kafka主题中获取数据并将数据输出到另一个Kafka主题.我写了一个非常基本的StreamTask但是在执行时我遇到了一个错误.

错误如下:

Exception in thread "main" org.apache.samza.SamzaException: org.apache.kafka.common.errors.TimeoutException: Failed to update metadata after 193 ms.
    at org.apache.samza.coordinator.stream.CoordinatorStreamSystemProducer.send(CoordinatorStreamSystemProducer.java:112)
    at org.apache.samza.coordinator.stream.CoordinatorStreamSystemProducer.writeConfig(CoordinatorStreamSystemProducer.java:129)
    at org.apache.samza.job.JobRunner.run(JobRunner.scala:79)
    at org.apache.samza.job.JobRunner$.main(JobRunner.scala:48)
    at org.apache.samza.job.JobRunner.main(JobRunner.scala)
 Caused by: org.apache.kafka.common.errors.TimeoutException: Failed to update metadata after 193 ms
Run Code Online (Sandbox Code Playgroud)

我不完全确定如何配置或让脚本编写所需的Kafka元数据.下面是我的StreamTask代码和属性文件.在属性文件中,我添加了元数据部分,以查看这是否有助于此后的过程,但无济于事.这是正确的方向还是我完全错过了什么?

import org.apache.samza.task.StreamTask;
import org.apache.samza.task.MessageCollector;
import org.apache.samza.task.TaskCoordinator;
import org.apache.samza.system.SystemStream;
import org.apache.samza.system.IncomingMessageEnvelope;
import org.apache.samza.system.OutgoingMessageEnvelope;

/*
*   Take all messages received and send them to
*   a Kafka topic called "words"
*/

public class TestStreamTask implements StreamTask{

    private static final SystemStream OUTPUT_STREAM = new SystemStream("kafka" , "words"); …
Run Code Online (Sandbox Code Playgroud)

java metadata apache-kafka apache-samza

8
推荐指数
1
解决办法
4857
查看次数

为什么YARN工作没有转换到RUNNING状态?

我有很多想要参加的Samza工作.我可以让第一个运行正常.但是,第二个作业似乎处于ACCEPTED状态,并且在我杀死第一个作业之前从未转换到RUNNING状态.

以下是YARN UI的视图:

YARN UI

以下是第二个作业的详细信息,您可以看到没有分配任何节点: 在此输入图像描述

我有2个数据节点,所以我应该可以运行多个作业.这是我的相关部分yarn-site.xml(我在文件中唯一的其他配置是与HA配置,Zookeeper等):

<property>
    <name>yarn.scheduler.minimum-allocation-mb</name>
    <value>128</value>
    <description>Minimum limit of memory to allocate to each container request at the Resource Manager.</description>
</property>
<property>
    <name>yarn.scheduler.maximum-allocation-mb</name>
    <value>2048</value>
    <description>Maximum limit of memory to allocate to each container request at the Resource Manager.</description>
</property>
<property>
    <name>yarn.scheduler.minimum-allocation-vcores</name>
    <value>1</value>
    <description>The minimum allocation for every container request at the RM, in terms of virtual CPU cores. Requests lower than this won't take effect, and the specified value will get allocated the minimum.</description>
</property>
<property> …
Run Code Online (Sandbox Code Playgroud)

hadoop hadoop-yarn apache-samza

8
推荐指数
1
解决办法
7029
查看次数

卡夫卡生产者超时异常

我正在运行将数据写入 Kafka 主题的 Samza 流作业。Kafka 正在运行一个 3 节点集群。Samza 作业部署在纱线上。我们在容器日志中看到了很多这些异常:

 INFO [2018-10-16 11:14:19,410] [U:2,151,F:455,T:2,606,M:2,658] samza.container.ContainerHeartbeatMonitor:[ContainerHeartbeatMonitor:stop:61] - [main] - Stopping ContainerHeartbeatMonitor
ERROR [2018-10-16 11:14:19,410] [U:2,151,F:455,T:2,606,M:2,658] samza.runtime.LocalContainerRunner:[LocalContainerRunner:run:107] - [main] - Container stopped with Exception. Exiting process now.
org.apache.samza.SamzaException: org.apache.samza.SamzaException: Unable to send message from TaskName-Partition 15 to system kafka.
        at org.apache.samza.task.AsyncRunLoop.run(AsyncRunLoop.java:147)
        at org.apache.samza.container.SamzaContainer.run(SamzaContainer.scala:694)
        at org.apache.samza.runtime.LocalContainerRunner.run(LocalContainerRunner.java:104)
        at org.apache.samza.runtime.LocalContainerRunner.main(LocalContainerRunner.java:149)
Caused by: org.apache.samza.SamzaException: Unable to send message from TaskName-Partition 15 to system kafka.
        at org.apache.samza.system.kafka.KafkaSystemProducer$$anon$1.onCompletion(KafkaSystemProducer.scala:181)
        at org.apache.kafka.clients.producer.internals.RecordBatch.done(RecordBatch.java:109)
        at org.apache.kafka.clients.producer.internals.RecordBatch.maybeExpire(RecordBatch.java:160)
        at org.apache.kafka.clients.producer.internals.RecordAccumulator.abortExpiredBatches(RecordAccumulator.java:245)
        at org.apache.kafka.clients.producer.internals.Sender.run(Sender.java:212)
        at org.apache.kafka.clients.producer.internals.Sender.run(Sender.java:135)
        at …
Run Code Online (Sandbox Code Playgroud)

apache-kafka apache-samza kafka-producer-api

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

Scala错误:未绑定的占位符参数和模式匹配条件

我正在尝试组合模式匹配和条件,但这段代码(这是一个Samza任务):

override def process(incomingMessageEnvelope: IncomingMessageEnvelope, messageCollector: MessageCollector, taskCoordinator: TaskCoordinator): Unit = {
    val event = (incomingMessageEnvelope getMessage).asInstanceOf[Map[String, Date]]
    val symbol = (event get "symbol").asInstanceOf[String]
    val eventDate = (event get "date").asInstanceOf[Date]

    (store get symbol) match {
      case x: java.util.Date if x.equals(eventDate) || x.after(eventDate) => _ 
      case _ => {
        this.store.put(symbol, eventDate)
      }
    }
  }
Run Code Online (Sandbox Code Playgroud)

返回此错误:

Error:(30, 38) unbound placeholder parameter
  case x if isGreaterOf(x, x) => _
                                 ^
Run Code Online (Sandbox Code Playgroud)

你知道错误吗?

谢谢

问候

赞布罗塔

scala pattern-matching apache-samza

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