【问题标题】:How to expose Java method for Kafka (commitSync with partitions) in Scala?如何在Scala中为Kafka(commitSync with partitions)公开Java方法?
【发布时间】:2023-03-03 04:29:01
【问题描述】:

我正在尝试通过 Scala 公开 Java 方法(有关原始 Java 方法的更多详细信息 - 它来自 Kafka

这是原始的 Java 方法:

public void commitSync(Map<TopicPartition,OffsetAndMetadata> offsets)

如何在 Scala 中向方法公开和传递参数?我有类似的东西:

def commitSync() = {
     consumer.commitSync(...)
}

谢谢。

【问题讨论】:

  • Scala sn-p 看起来是正确的,但是您可能希望将集合从 Java 的集合转换为 Scala 的集合,例如 java.util.Mapscala.collection.immutable.Map。但是,我不确定您所说的回调 Scala 中的方法是什么意思。您能否详细说明您想在上面的示例中传递什么样的回调以及传递给哪个方法?
  • 你能分享一个完整的代码示例吗? (我还想确定是否有东西应该放在点内......)。我是 Java/Scala 世界的新手..(已编辑问题,谢谢)。

标签: java scala apache-kafka


【解决方案1】:

你的 Scala sn-p 看起来是正确的,这就是我将如何填充它的其余部分:

import org.apache.kafka.clients.consumer.{KafkaConsumer, OffsetAndMetadata}
import org.apache.kafka.common.TopicPartition

import collection.mutable.Map
import collection.JavaConverters._

//initialise your consumer the way you want
val consumer = createKafkaConsumer(config, subscriptions)

//you could accept a scala.collection.mutable.Map here
def commitSync(offsets: Map[TopicPartition, OffsetAndMetadata]) = {
    //and then convert it to a java.util.Map
    consumer.commitSync(offsets.asJava)
}

【讨论】:

  • 谢谢!什么是.asJava 关键字?那有必要吗?我也看到 val consumer = new KafkaConsumer[K,V](Map.empty.asJava) ,我目前拥有它为 var consumer = createKafkaConsumer(config, subscriptions).. 对此有什么想法吗?
  • 啊,对不起!我应该更明确一点。 .asJava 是由上面的 JavaConversions 导入提供的隐式装饰器。至于valvar,当你想要一个不可变的值引用时,你会使用val。即:您以后将无法将其重新分配给其他任何东西。就使用createKafkaConsumer(config, subscriptions) 而言,没关系。我将修改我的答案以包括这一点。刚写答案的时候,我只是没想太多new KafkaConsumer[K,V](Map.empty.asJava)
  • 太棒了!你可以做我的导师吗? :)
  • 还有一个问题 - 您是否有示例说明如何“调用”此方法 def commitSync(带参数)。我正在尝试为它做一些测试用例..
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-01-30
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-11-11
相关资源
最近更新 更多