【发布时间】:2017-08-28 04:08:16
【问题描述】:
我只是想知道是否可以将使用 Java 的 Executor 服务编写的旧多线程代码替换为 Akka。我对此毫无疑问。
Is akka actor runs in their own thread?
How Threads will be assigned for the Actors ?
What are the pros and cons of migration of it is possible?
目前我使用固定线程池进行多线程,并提交一个可调用的。
示例代码,
public class KafkaConsumerFactory {
private static Map<String,KafkaConsumer> registry = new HashMap<>();
private static ThreadLocal<KafkaConsumer> consumers = new ThreadLocal<KafkaConsumer>(){
@Override
protected KafkaConsumer initialValue() {
return new KafkaConsumer(createConsumerConfig());
}
};
static {
Runtime.getRuntime().addShutdownHook(new Thread(){
@Override
public void run() {
registry.forEach((tid,con) -> {
try{
con.close();
} finally {
System.out.println("Yes!! Consumer for " + tid + " is closed.");
}
});
}
});
}
private static Properties createConsumerConfig() {
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "newcon-grp5");
props.put("key.deserializer", StringDeserializer.class.getName());
props.put("value.deserializer", KafkaKryoSerde.class.getName());
return props;
}
public static <K,V> KafkaConsumer<K,V> createConsumer(){
registry.put(Thread.currentThread().getName(),consumers.get());
return consumers.get();
}
}
/////////////////////////////////////// //////////
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.util.*;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
public class KafkaNewConsumer {
public static int MAX_THREADS = 10;
private ExecutorService es = null;
private boolean stopRequest = false;
public static void main(String[] args){
KafkaNewConsumer knc = new KafkaNewConsumer();
Runtime.getRuntime().addShutdownHook(new Thread(){
@Override
public void run(){
knc.es.shutdown();
try {
knc.es.awaitTermination(500, TimeUnit.MILLISECONDS);
} catch (InterruptedException ignored) {
}finally {
System.out.println("Finished");
}
}
});
knc.consumeTopic("rtest3",knc::recordConsuemer);
}
public void recordConsuemer(ConsumerRecord<?,?> record){
String result = new StringJoiner(": ")
.add(Thread.currentThread().getName())
.add("ts").add(String.valueOf(record.timestamp()))
.add("offset").add(String.valueOf(record.offset()))
.add("data").add(String.valueOf(record.value()))
.add("value-len").add(String.valueOf(record.serializedValueSize()))
.toString();
System.out.println(result);
}
public void consumeTopic(String topicName, Consumer<ConsumerRecord<?,?>> fun){
KafkaConsumer con= KafkaConsumerFactory.createConsumer();
int paritions = con.partitionsFor(topicName).size();
int noOfThread = (MAX_THREADS < paritions) ? MAX_THREADS :paritions;
es = Executors.newFixedThreadPool(noOfThread);
con.close();
for(int i=0;i<noOfThread;i++){
es.submit(()->{
KafkaConsumer consumer = KafkaConsumerFactory.createConsumer();
try{
while (!stopRequest){
consumer.subscribe(Collections.singletonList(topicName));
ConsumerRecords<?,?> records = consumer.poll(5000);
records.forEach(fun);
consumer.commitSync();
}
}catch(Exception e){
e.printStackTrace();
} finally {
consumer.close();
}
});
}
}
}
上网查了一些教程,有的直接总结
actors 非常好,比传统线程更快。
但没有解释它如何变得比线程更快?
我尝试了一些示例 Akka(Akka sample from activator) 代码,并在所有参与者中打印了 Thread.currentThread.getName 并发现创建了名为 (helloakka-akka.actor.default-dispatcher-X) 的不同调度程序线程。
但是怎么做呢?谁在创建这些线程?他们的配置在哪里?线程和Actor的映射关系是什么?
每次我发送消息时,Akka 都会创建新线程吗?还是内部使用了线程池?
如果我需要 100 个线程来并行执行某些任务的某些部分,我是否需要创建 100 个 Actor 并向每个 Actor 发送 1 条消息?或者我需要创建 1 个演员并将 100 条消息放入它的队列中,它将被分叉成 100 个线程。
真的很困惑
【问题讨论】:
-
我可以知道投反对票的原因吗?
-
信息不足。什么是“代码”?您要更换哪种线程模型,为什么?你如何同步?你把结果推到哪里?什么是阻塞,什么不是?前后真相的来源在哪里?现在至少回答其中的几个问题可能会变成一个真正的问题 - 我看不出那里甚至可以回答什么。因为如果你问“我可以吗” - 唯一的答案是“是的,你技术上可以”。
-
一个简单的 System out println 对我来说就足够了,我只是对 Actor 将在其中运行的线程的映射以及如何管理这些线程感到困惑。我修改了我的问题。任何指针都会有所帮助。
-
@M.Prokhorov 示例代码也已添加。
标签: java multithreading akka