Tiny并行计算框架之复杂示例

问题来源

非常感谢@doctorwho的问题:

假如职业介绍所来了一批生产汽车的工作,假设生产一辆汽车任务是这样的:搭好底盘、拧4个轮胎、安装发动机、安装4个座椅、再装4个车门、最后安装顶棚。之间有的任务是可以并行计算的(比如拧4个轮胎,安装发动机和安装座椅),有的任务有前置任务(比如先装好座椅,才能装车门和顶棚)。让两组包工头组织两种类型的工作:将工人分成两种类型,即可并行计算的放在同一组内,由职业介绍所来控制A组包工头做完的任务交给B组包工头。中间环节的半成品保存到Warehouse中,是这样使用TINY框架来生产汽车么?

接下来,我就用Tiny并行计算框架来展示一下这个示例,在编写示例的时候,发现了一个BUG,这也充分体现了开源的精神与价值,再次感谢@doctorwho。

问题分析

doctorwho的问题还是比较复杂的,但是实际上道理是一样的,因此我把问题简化成下面的过程

第一步:构建底盘

第二步:并行进行安装引擎,座位和轮胎

第三步:并行进行安装门及车顶

由于我和doctorwho都不是造车行家,因此就不用纠结这么造是不是合理了,假设这么做就是合理的。

代码实现

按我前面说的过程,工人是必须要有的,因此我们首先构建工人:

第一步的底盘构建工人

1public class StepFirstWorker extends AbstractWorker { 2 public StepFirstWorker() throws RemoteException { 3 super("first"); 4 } 5 6 @Override 7 protected Warehouse doWork(Work work) throws RemoteException { 8 System.out.println(String.format("%s 构建底盘完成.", work.getInputWarehouse().get("carType"))); 9 Warehouse outputWarehouse = work.getInputWarehouse(); 10 outputWarehouse.put("baseInfo", "something about baseInfo"); 11 return outputWarehouse; 12 } 13}

由于第二步工人有好几个类型,因此再搞个第二步抽象工人:

1public abstract class StepThirdWorker extends AbstractWorker { 2    public StepThirdWorker() throws RemoteException { 3        super("third"); 4    } 5 6 7    protected boolean acceptMyWork(Work work) { 8        String workClass = work.getInputWarehouse().get("class"); 9        if (workClass != null) { 10            return true; 11        } 12        return false; 13    } 14    protected Warehouse doMyWork(Work work) throws RemoteException { 15        System.out.println(String.format("Base:%s ", work.getInputWarehouse().get("baseInfo"))); 16        System.out.println(String.format("%s is Ok", work.getInputWarehouse().get("class"))); 17        return work.getInputWarehouse(); 18    } 19}

接下来构建第二步的引擎工人:

1public class StepSecondEngineWorker extends StepSecondWorker { 2 3 4 public static final String ENGINE = "engine"; 5 6 7 public StepSecondEngineWorker() throws RemoteException { 8 super(); 9 } 10 11 12 public boolean acceptWork(Work work) { 13 return acceptMyWork(work); 14 } 15 16 17 protected Warehouse doWork(Work work) throws RemoteException { 18 return super.doMyWork(work); 19 } 20}

第二步的座位工人:

1public class StepSecondSeatWorker extends StepSecondWorker { 2 3 4    public static final String SEAT = "seat"; 5 6 7    public StepSecondSeatWorker() throws RemoteException { 8        super(); 9    } 10    public boolean acceptWork(Work work) { 11       return acceptMyWork(work); 12    } 13    protected Warehouse doWork(Work work) throws RemoteException { 14        return super.doMyWork(work); 15    } 16} 17 18 19 20 21 22 23 24 25 26 27

第二步的轮胎工人:

1public class StepSecondTyreWorker extends StepSecondWorker { 2 public static final String TYRE = "tyre"; 3 4 5 public StepSecondTyreWorker() throws RemoteException { 6 super(); 7 } 8 9 10 public boolean acceptWork(Work work) { 11 return acceptMyWork(work); 12 } 13 14 15 protected Warehouse doWork(Work work) throws RemoteException { 16 return super.doMyWork(work); 17 } 18}

同理,第三步也是大同小异的。

第三步的抽象工人类:

1public abstract class StepThirdWorker extends AbstractWorker { 2    public StepThirdWorker() throws RemoteException { 3        super("third"); 4    } 5 6 7    protected boolean acceptMyWork(Work work) { 8        String workClass = work.getInputWarehouse().get("class"); 9        if (workClass != null) { 10            return true; 11        } 12        return false; 13    } 14    protected Warehouse doMyWork(Work work) throws RemoteException { 15        System.out.println(String.format("Base:%s ", work.getInputWarehouse().get("baseInfo"))); 16        System.out.println(String.format("%s is Ok", work.getInputWarehouse().get("class"))); 17        return work.getInputWarehouse(); 18    } 19}

第三步的车门工人:

1public class StepThirdDoorWorker extends StepThirdWorker { 2 3 4    public static final String DOOR = "door"; 5 6 7    public StepThirdDoorWorker() throws RemoteException { 8        super(); 9    } 10    public boolean acceptWork(Work work) { 11        return acceptMyWork(work); 12    } 13    @Override 14    protected Warehouse doWork(Work work) throws RemoteException { 15        return super.doMyWork(work); 16    } 17}

第三步的车顶工人:

1public class StepThirdRoofWorker extends StepThirdWorker { 2 3 4    public static final String ROOF = "roof"; 5 6 7    public StepThirdRoofWorker() throws RemoteException { 8        super(); 9    } 10    public boolean acceptWork(Work work) { 11        return acceptMyWork(work); 12    } 13    protected Warehouse doWork(Work work) throws RemoteException { 14        return super.doMyWork(work); 15    } 16}

以上就把工人都构建好了,我们前面也说过,如果要进行任务分解,是必须要构建任务分解合并器的,这里简单起见,只实现任务分解了。

第二部的任务分解:

1public class SecondWorkSplitter implements WorkSplitter { 2 public List<Warehouse> split(Work work, List<Worker> workers) throws RemoteException { 3 List<Warehouse> list = new ArrayList<Warehouse>(); 4 list.add(getWareHouse(work.getInputWarehouse(), "engine")); 5 list.add(getWareHouse(work.getInputWarehouse(), "seat")); 6 list.add(getWareHouse(work.getInputWarehouse(), "seat")); 7 list.add(getWareHouse(work.getInputWarehouse(), "seat")); 8 list.add(getWareHouse(work.getInputWarehouse(), "seat")); 9 list.add(getWareHouse(work.getInputWarehouse(), "tyre")); 10 list.add(getWareHouse(work.getInputWarehouse(), "tyre")); 11 list.add(getWareHouse(work.getInputWarehouse(), "tyre")); 12 list.add(getWareHouse(work.getInputWarehouse(), "tyre")); 13 return list; 14 } 15 16 private Warehouse getWareHouse(Warehouse inputWarehouse, String stepClass) { 17 Warehouse warehouse = new WarehouseDefault(); 18 warehouse.put("class", stepClass); 19 warehouse.putSubWarehouse(inputWarehouse); 20 return warehouse; 21 } 22}

从上面可以看到,构建了一个引擎的仓库,4个座位仓库,4个轮胎仓库。呵呵,既然能并行,为啥不让他做得更快些?

接下来是第三步的任务分解器:

1public class ThirdWorkSplitter implements WorkSplitter { 2 public List<Warehouse> split(Work work, List<Worker> workers) throws RemoteException { 3 List<Warehouse> list = new ArrayList<Warehouse>(); 4 list.add(getWareHouse(work.getInputWarehouse(), "door")); 5 list.add(getWareHouse(work.getInputWarehouse(), "door")); 6 list.add(getWareHouse(work.getInputWarehouse(), "door")); 7 list.add(getWareHouse(work.getInputWarehouse(), "door")); 8 list.add(getWareHouse(work.getInputWarehouse(), "roof")); 9 return list; 10 } 11 12 private Warehouse getWareHouse(Warehouse inputWarehouse, String stepClass) { 13 Warehouse warehouse = new WarehouseDefault(); 14 warehouse.put("class", stepClass); 15 warehouse.putSubWarehouse(inputWarehouse); 16 return warehouse; 17 } 18}

从上面可以看到,第三部构建了4个门仓库一个车顶仓库,同样的,可以让4个工人同时装门。

上面就把所有的准备工作都做好了,接下来就是测试方法了:

1public class Test { 2 public static void main(String[] args) throws IOException, ClassNotFoundException, InterruptedException { 3 JobCenter jobCenter = new JobCenterLocal(); 4 5 6 for (int i = 0; i < 5; i++) { 7 jobCenter.registerWorker(new StepFirstWorker()); 8 } 9 for (int i = 0; i < 5; i++) { 10 jobCenter.registerWorker(new StepSecondTyreWorker()); 11 } 12 for (int i = 0; i < 5; i++) { 13 jobCenter.registerWorker(new StepSecondSeatWorker()); 14 } 15 for (int i = 0; i < 5; i++) { 16 jobCenter.registerWorker(new StepSecondEngineWorker()); 17 } 18 for (int i = 0; i < 5; i++) { 19 jobCenter.registerWorker(new StepThirdDoorWorker()); 20 } 21 for (int i = 0; i < 5; i++) { 22 jobCenter.registerWorker(new StepThirdRoofWorker()); 23 } 24 25 jobCenter.registerForeman(new ForemanSelectOneWorker("first")); 26      jobCenter.registerForeman(new ForemanSelectAllWorker("second", 27 new SecondWorkSplitter())); 28 jobCenter.registerForeman(new ForemanSelectAllWorker("third",new ThirdWorkSplitter())); 29 30 31 Warehouse inputWarehouse = new WarehouseDefault(); 32 inputWarehouse.put("class", "car"); 33 inputWarehouse.put("carType", "普桑"); 34 WorkDefault work = new WorkDefault("first", inputWarehouse); 35 work.setForemanType("first"); 36 WorkDefault work2 = new WorkDefault("second"); 37 work2.setForemanType("second"); 38 WorkDefault work3 = new WorkDefault("third"); 39 work3.setForemanType("third"); 40 work.setNextWork(work2).setNextWork(work3); 41 42 43 Warehouse warehouse = jobCenter.doWork(work); 44 45 46 jobCenter.stop(); 47 48 } 49}

呵呵,工人各加了5个,然后注册了三个工头,第一步的工头是随便挑一个工人类型的,第二步和第三步是挑所有工人的,同时还指定了任务分解器。

接下来就构建了一个工作,造一个高端大气上档次的普桑汽车,然后告诉职业介绍所说给我造就可以了。

下面是造车的过程,我把日志也贴上来了:

1普桑 构建底盘完成. 2-234 [RMI TCP Connection(1)-192.168.84.73] INFO - 线程组<id:4af96b81d14a4954a6b649308d444e4c,type:second>运行开始,线程数9... 3-234 [id:4af96b81d14a4954a6b649308d444e4c,type:second-a763f156ffd74b5db285198d2498edcf] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-a763f156ffd74b5db285198d2498edcf>运行开始... 4-234 [id:4af96b81d14a4954a6b649308d444e4c,type:second-c2d2fb38ef6c4509b3a39b3e7d5c1d61] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-c2d2fb38ef6c4509b3a39b3e7d5c1d61>运行开始... 5-234 [id:4af96b81d14a4954a6b649308d444e4c,type:second-d624ea0a6df3409c80df6b97ab3c813b] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-d624ea0a6df3409c80df6b97ab3c813b>运行开始... 6-235 [id:4af96b81d14a4954a6b649308d444e4c,type:second-abdb57f0641a4727a9efa744d07cf2d1] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-abdb57f0641a4727a9efa744d07cf2d1>运行开始... 7-236 [id:4af96b81d14a4954a6b649308d444e4c,type:second-d6f7074f6c4a4b12bd37ec5f5c11aff8] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-d6f7074f6c4a4b12bd37ec5f5c11aff8>运行开始... 8-237 [id:4af96b81d14a4954a6b649308d444e4c,type:second-04db3f945b804500a2bbe2b9aabdce3b] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-04db3f945b804500a2bbe2b9aabdce3b>运行开始... 9Base:something about baseInfo 10engine is Ok 11Base:something about baseInfo 12seat is Ok 13-245 [id:4af96b81d14a4954a6b649308d444e4c,type:second-a763f156ffd74b5db285198d2498edcf] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-a763f156ffd74b5db285198d2498edcf>运行结束 14-246 [id:4af96b81d14a4954a6b649308d444e4c,type:second-f3efba2dc7804c6cbcd5a25f42fdc177] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-f3efba2dc7804c6cbcd5a25f42fdc177>运行开始... 15-246 [id:4af96b81d14a4954a6b649308d444e4c,type:second-abdb57f0641a4727a9efa744d07cf2d1] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-abdb57f0641a4727a9efa744d07cf2d1>运行结束 16Base:something about baseInfo 17seat is Ok 18-248 [id:4af96b81d14a4954a6b649308d444e4c,type:second-8c3b9359bcfa4de7b6e0492daab0d73a] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-8c3b9359bcfa4de7b6e0492daab0d73a>运行开始... 19Base:something about baseInfo 20tyre is Ok 21-250 [id:4af96b81d14a4954a6b649308d444e4c,type:second-f3efba2dc7804c6cbcd5a25f42fdc177] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-f3efba2dc7804c6cbcd5a25f42fdc177>运行结束 22-250 [id:4af96b81d14a4954a6b649308d444e4c,type:second-c2d2fb38ef6c4509b3a39b3e7d5c1d61] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-c2d2fb38ef6c4509b3a39b3e7d5c1d61>运行结束 23Base:something about baseInfo 24seat is Ok 25-252 [id:4af96b81d14a4954a6b649308d444e4c,type:second-869b573e226046aca8ad30765f1f300c] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-869b573e226046aca8ad30765f1f300c>运行开始... 26-253 [id:4af96b81d14a4954a6b649308d444e4c,type:second-8c3b9359bcfa4de7b6e0492daab0d73a] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-8c3b9359bcfa4de7b6e0492daab0d73a>运行结束 27Base:something about baseInfo 28seat is Ok 29Base:something about baseInfo 30tyre is Ok 31-257 [id:4af96b81d14a4954a6b649308d444e4c,type:second-d624ea0a6df3409c80df6b97ab3c813b] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-d624ea0a6df3409c80df6b97ab3c813b>运行结束 32-258 [id:4af96b81d14a4954a6b649308d444e4c,type:second-04db3f945b804500a2bbe2b9aabdce3b] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-04db3f945b804500a2bbe2b9aabdce3b>运行结束 33Base:something about baseInfo 34tyre is Ok 35Base:something about baseInfo 36tyre is Ok 37-262 [id:4af96b81d14a4954a6b649308d444e4c,type:second-869b573e226046aca8ad30765f1f300c] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-869b573e226046aca8ad30765f1f300c>运行结束 38-264 [id:4af96b81d14a4954a6b649308d444e4c,type:second-d6f7074f6c4a4b12bd37ec5f5c11aff8] INFO - 线程<id:4af96b81d14a4954a6b649308d444e4c,type:second-d6f7074f6c4a4b12bd37ec5f5c11aff8>运行结束 39-264 [RMI TCP Connection(1)-192.168.84.73] INFO - 线程组<id:4af96b81d14a4954a6b649308d444e4c,type:second>运行结束, 用时:30ms 40-333 [RMI TCP Connection(1)-192.168.84.73] INFO - 线程组<id:9b0de678632d4fb8b87ae9db4b6436f8,type:third>运行开始,线程数5... 41-334 [id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-ca58e9b733514c668a224875417c9d26] INFO - 线程<id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-ca58e9b733514c668a224875417c9d26>运行开始... 42-334 [id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-17debef817e34c49996a2c38840f3de2] INFO - 线程<id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-17debef817e34c49996a2c38840f3de2>运行开始... 43-334 [id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-af0e914cce89480987c6184a885770d5] INFO - 线程<id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-af0e914cce89480987c6184a885770d5>运行开始... 44-334 [id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-3f7707cb2a224d3d8844a09271b24a07] INFO - 线程<id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-3f7707cb2a224d3d8844a09271b24a07>运行开始... 45Base:something about baseInfo 46door is Ok 47Base:something about baseInfo 48door is Ok 49Base:something about baseInfo 50-338 [id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-ca58e9b733514c668a224875417c9d26] INFO - 线程<id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-ca58e9b733514c668a224875417c9d26>运行结束 51door is Ok 52-339 [id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-cfccfb37ebda4279b7552f7e060b2ddb] INFO - 线程<id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-cfccfb37ebda4279b7552f7e060b2ddb>运行开始... 53-338 [id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-17debef817e34c49996a2c38840f3de2] INFO - 线程<id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-17debef817e34c49996a2c38840f3de2>运行结束 54Base:something about baseInfo 55door is Ok 56Base:something about baseInfo 57-340 [id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-3f7707cb2a224d3d8844a09271b24a07] INFO - 线程<id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-3f7707cb2a224d3d8844a09271b24a07>运行结束 58roof is Ok 59-340 [id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-af0e914cce89480987c6184a885770d5] INFO - 线程<id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-af0e914cce89480987c6184a885770d5>运行结束 60-342 [id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-cfccfb37ebda4279b7552f7e060b2ddb] INFO - 线程<id:9b0de678632d4fb8b87ae9db4b6436f8,type:third-cfccfb37ebda4279b7552f7e060b2ddb>运行结束 61-343 [RMI TCP Connection(1)-192.168.84.73] INFO - 线程组<id:9b0de678632d4fb8b87ae9db4b6436f8,type:third>运行结束, 用时:10ms

从上面的日志可以看出:

由于第一步工作是挑一个单干的,因此是没有启用线程组的

第二步同时有9个线程干活:

1-234 [RMI TCP Connection(1)-192.168.84.73] INFO - 线程组<id:4af96b81d14a4954a6b649308d444e4c,type:second>运行开始,线程数9... 2... 3-264 [RMI TCP Connection(1)-192.168.84.73] INFO - 线程组<id:4af96b81d14a4954a6b649308d444e4c,type:second>运行结束, 用时:30ms

第三步同时有5个线程干活:

1-333 [RMI TCP Connection(1)-192.168.84.73] INFO - 线程组<id:9b0de678632d4fb8b87ae9db4b6436f8,type:third>运行开始,线程数5... 2... 3-343 [RMI TCP Connection(1)-192.168.84.73] INFO - 线程组<id:9b0de678632d4fb8b87ae9db4b6436f8,type:third>运行结束, 用时:10ms

总结:

Tiny并行计算框架确实是可以方便的解决各种复杂并行计算的问题。

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )