【问题标题】:Migrate Java multi threading with Executor Service to Akka使用 Executor Service 将 Java 多线程迁移到 Akka
【发布时间】: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


【解决方案1】:

迁移到 Actor 系统对于基于 executor 的系统来说不是一项小任务,但它可以完成。它要求您重新思考设计系统的方式并考虑参与者的影响。例如,在线程架构中,您为业务流程创建一些处理程序,将其放入可运行文件中,然后让它在线程上执行操作。这对于演员范式来说是完全不合适的。您必须重新构建系统以处理消息传递并使用消息传递来调用任务。此外,您还必须将您对业务流程的思考方式从命令式方法转变为基于消息的方法。以购买产品的简单任务为例。我假设您知道如何在执行程序中执行此操作。在演员系统中,您可以这样做:

(Purchase Product) -> UserActor -> (BillCredit Card) -> CCProcessing Actor -> (Purchase Approved and Billed Item) -> inventory manager -> ... 等等

在每个阶段,括号中的内容是发送给相关参与者的异步消息,该参与者执行业务逻辑,然后将消息转发给流程中的下一个参与者。

现在这只是创建基于actor的系统的一种方法,还有许多其他技术,但核心基础是您不能强制思考,而是作为每个独立运行的步骤的集合。然后消息以常规顺序在系统中爆炸,但您无法确定顺序,或者即使消息会到达那里,所以您必须设计语义来处理它。在上面的系统中,我可能有另一个参与者每两分钟检查一次尚未提交给计费的孤立订单。当然,这意味着我的消息需要具有幂等性,以确保如果我第二次发送它们就可以了,它们不会向用户收费两次。

我知道我没有处理您的具体示例,我只是想为您提供一些上下文,即演员不仅仅是创建执行者的另一种方式(我想您可以那样滥用它们,但这是不可取的)而是完全不同的设计范式。一个非常值得学习的范式,如果你实现了飞跃,你将永远不想再做 executors。

【讨论】:

  • 感谢您的回复。我仍然有点困惑。我尝试了一些示例 Akka(Akka sample from activator)代码,并在所有参与者中打印了 Thread.currentThread.getName 并发现名为(helloakka-akka.actor.default-dispatcher-X)的不同调度程序线程是创建的。但是如何?谁在创建这些线程?他们的配置在哪里?线程和Actor的映射关系是什么?
  • 每次我发送消息时Akka都会创建新线程吗?如果我需要 100 个线程,是否需要为每个线程创建 100 个 Actor 和 1 条消息?或者我需要创建 1 个演员并将 100 条消息放入它的队列中,它将被分叉成 100 个线程。真的很困惑。
  • 确定关于演员的一件事。不要再考虑线程了。不,真的..停下来! :) 我看到你在想他们。当您收到一条消息时,您将其发送给参与者,让调度程序为您管理线程。在那个演员里面,你是在一个快乐的环境中。一次只能处理一条消息。您无需担心干扰。如果有公共资源,请使用他们自己的演员对这些资源进行建模,向他们发送一条消息,他们会在接到您的电话时向您发送一条消息,用他们添加的额外信息来装饰您的请求。
  • 那么你所要做的就是确保你有足够的演员来处理带宽。在路由器中部署池?如果消息一次是每个节点一个,则它是一个普通的参与者。如果您有一个进程必须在整个集群上一次只发生一个,请将参与者集群部署为分片。尽快将您正在执行的任务保留在消息中。请记住,演员几乎不应该等待,否则您将成为应用程序的瓶颈。将任何长期任务放在 future 中并学习 abotu become() 来实现状态机。我知道要吞下很多东西,只需从一个简单的演员开始,然后继续努力。
猜你喜欢
  • 2019-03-21
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-01-07
  • 2011-12-21
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多