我目前正在尝试在php中实现一个作业队列.然后,队列将作为批处理作业处理,并且应该能够并行处理某些作业.
我已经做了一些研究并找到了几种方法来实现它,但我并没有真正意识到它们的优点和缺点.
例如,通过多次调用脚本来执行并行处理,fsockopen如下所述:
PHP中的简单并行处理
我找到的另一种方法是使用这些curl_multi功能.
curl_multi_exec PHP文档
但我认为这两种方式会增加相当多的开销,在队列上创建批处理应该主要在后台运行吗?
我也读到了关于pcntl_fork这似乎也是一种处理问题的方法.但是如果你真的不知道自己在做什么(就像我现在这样;),那看起来就会变得非常混乱.)
我也看了一下Gearman,但在那里我还需要根据需要动态生成工作线程,而不是只运行一些,然后让gearman作业服务器将它发送给自由工作者.特别是因为线程应该在执行一个作业后干净利落地退出,不会遇到最终的内存泄漏(在该问题中代码可能不完美).
Gearman入门
所以我的问题是,你如何处理PHP中的并行处理?为什么选择你的方法,不同的方法有哪些优点/缺点?
感谢您的任何意见.
我想更好地理解我在C#中的Async和Parallel选项.在下面的片段中,我列出了我最常遇到的5种方法.但我不确定选择哪个 - 或者更好的是,在选择时要考虑的标准:
方法1:任务
(见http://msdn.microsoft.com/en-us/library/dd321439.aspx)
调用StartNew在功能上等同于使用其构造函数之一创建Task,然后调用Start来安排执行.但是,除非必须分离创建和调度,否则StartNew是简化和性能的推荐方法.
TaskFactory的StartNew方法应该是创建和调度计算任务的首选机制,但是对于必须分离创建和调度的场景,可以使用构造函数,然后可以使用任务的Start方法来安排任务以便稍后执行时间.
// using System.Threading.Tasks.Task.Factory
void Do_1()
{
var _List = GetList();
_List.ForEach(i => Task.Factory.StartNew(_ => { DoSomething(i); }));
}
Run Code Online (Sandbox Code Playgroud)
方法2:QueueUserWorkItem
(参见http://msdn.microsoft.com/en-us/library/system.threading.threadpool.getmaxthreads.aspx)
您可以将系统内存允许的线程池请求排队.如果请求多于线程池线程,则其他请求将保持排队,直到线程池线程可用.
您可以将排队方法所需的数据放在定义方法的类的实例字段中,也可以使用接受包含必要数据的对象的QueueUserWorkItem(WaitCallback,Object)重载.
// using System.Threading.ThreadPool
void Do_2()
{
var _List = GetList();
var _Action = new WaitCallback((o) => { DoSomething(o); });
_List.ForEach(x => ThreadPool.QueueUserWorkItem(_Action));
}
Run Code Online (Sandbox Code Playgroud)
方法3:Parallel.Foreach
(参见:http://msdn.microsoft.com/en-us/library/system.threading.tasks.parallel.foreach.aspx)
Parallel类为常见操作提供基于库的数据并行替换,例如for循环,每个循环以及一组语句的执行.
为可枚举源中的每个元素调用一次body委托.它以当前元素作为参数提供.
// using System.Threading.Tasks.Parallel
void Do_3()
{
var _List = GetList();
var _Action = new Action<object>((o) => …Run Code Online (Sandbox Code Playgroud) Java 5引入了Executor框架形式的线程池对异步任务执行的支持,其核心是java.util.concurrent.ThreadPoolExecutor实现的线程池.Java 7以java.util.concurrent.ForkJoinPool的形式添加了一个备用线程池.
查看各自的API,ForkJoinPool在标准场景中提供了ThreadPoolExecutor功能的超集(虽然严格来说ThreadPoolExecutor提供了比ForkJoinPool更多的调优机会).除此之外,fork/join任务看起来更快(可能是因为工作窃取调度程序)的观察结果显然需要更少的线程(由于非阻塞连接操作),可能会让人觉得ThreadPoolExecutor已被取代ForkJoinPool.
但这真的是对的吗?我读过的所有材料似乎总结为两种类型的线程池之间相当模糊的区别:
这种区别是否正确?我们能说出更具体的内容吗?
java parallel-processing threadpool threadpoolexecutor forkjoinpool
我在这里问了一个相关的问题并且响应运行良好: 使用parallel的parLapply:无法访问并行代码中的变量
问题是,当我尝试使用函数内部的答案时,它将无法工作,因为我认为它具有默认环境clusterExport.我已经阅读了小插图并查看了帮助文件,但我的知识库非常有限.我使用的方式我parLapply期望它的行为类似lapply但似乎没有.
这是我的尝试:
par.test <- function(text.var, gc.rate=10){
ntv <- length(text.var)
require(parallel)
pos <- function(i) {
paste(sapply(strsplit(tolower(i), " "), nchar), collapse=" | ")
}
cl <- makeCluster(mc <- getOption("cl.cores", 4))
clusterExport(cl=cl, varlist=c("text.var", "ntv", "gc.rate", "pos"))
parLapply(cl, seq_len(ntv), function(i) {
x <- pos(text.var[i])
if (i%%gc.rate==0) gc()
return(x)
}
)
}
par.test(rep("I like cake and ice cream so much!", 20))
#gives this error message
> par.test(rep("I like cake and ice cream so much!", 20))
Error in …Run Code Online (Sandbox Code Playgroud) 如果我在foreach... %dopar%没有注册集群的情况下运行,foreach会发出警告,并按顺序执行代码:
library("doParallel")
foreach(i=1:3) %dopar%
sqrt(i)
Run Code Online (Sandbox Code Playgroud)
产量:
Warning message:
executing %dopar% sequentially: no parallel backend registered
Run Code Online (Sandbox Code Playgroud)
但是,如果我在启动,注册和停止集群后运行相同的代码,则会失败:
cl <- makeCluster(2)
registerDoParallel(cl)
stopCluster(cl)
rm(cl)
foreach(i=1:3) %dopar%
sqrt(i)
Run Code Online (Sandbox Code Playgroud)
产量:
Error in summary.connection(connection) : invalid connection
Run Code Online (Sandbox Code Playgroud)
有没有相反的registerDoParallel()清理群集注册?还是我坚持使用旧集群的鬼魂,直到我重新开始我的R会话?
/编辑:一些谷歌搜索揭示bumphunter:::foreachCleanup()了bumphunter Biocondoctor包中的功能:
function ()
{
if (exists(".revoDoParCluster", where = doParallel:::.options)) {
if (!is.null(doParallel:::.options$.revoDoParCluster))
stopCluster(doParallel:::.options$.revoDoParCluster)
remove(".revoDoParCluster", envir = doParallel:::.options)
}
}
<environment: namespace:bumphunter>
Run Code Online (Sandbox Code Playgroud)
但是,此功能似乎无法解决问题.
library(bumphunter)
cl <- makeCluster(2)
registerDoParallel(cl)
stopCluster(cl)
rm(cl)
bumphunter:::foreachCleanup()
foreach(i=1:3) %dopar%
sqrt(i)
Run Code Online (Sandbox Code Playgroud)
foreach在哪里保留注册集群的信息?
环境:Ubuntu x86_64(14.10),Oracle JDK 1.8u25
我尝试使用并行流Files.lines()但我想要.skip()第一行(它是带有标题的CSV文件).所以我试着这样做:
try (
final Stream<String> stream = Files.lines(thePath, StandardCharsets.UTF_8)
.skip(1L).parallel();
) {
// etc
}
Run Code Online (Sandbox Code Playgroud)
但是后来一列未能解析成一个int ...
所以我尝试了一些简单的代码.文件问题很简单:
$ cat info.csv
startDate;treeDepth;nrMatchers;nrLines;nrChars;nrCodePoints;nrNodes
1422758875023;34;54;151;4375;4375;27486
$
Run Code Online (Sandbox Code Playgroud)
代码同样简单:
public static void main(final String... args)
{
final Path path = Paths.get("/home/fge/tmp/dd/info.csv");
Files.lines(path, StandardCharsets.UTF_8).skip(1L).parallel()
.forEach(System.out::println);
}
Run Code Online (Sandbox Code Playgroud)
我系统地得到以下结果(好吧,我只运行了大约20次):
startDate;treeDepth;nrMatchers;nrLines;nrChars;nrCodePoints;nrNodes
Run Code Online (Sandbox Code Playgroud)
我在这里错过了什么?
编辑似乎问题或误解比这更根深蒂固(下面的两个例子是由FreeNode的## java编写的):
public static void main(final String... args)
{
new BufferedReader(new StringReader("Hello\nWorld")).lines()
.skip(1L).parallel()
.forEach(System.out::println);
final Iterator<String> iter
= Arrays.asList("Hello", "World").iterator();
final Spliterator<String> spliterator
= Spliterators.spliteratorUnknownSize(iter, …Run Code Online (Sandbox Code Playgroud) 以下伪代码是否是线程安全的?
IList<T> dataList = SomeNhibernateRepository.GetData();
Parallel.For(..i..)
{
foreach(var item in dataList)
{
DoSomething(item);
}
}
Run Code Online (Sandbox Code Playgroud)
列表永远不会改变,它只是迭代并且并行读取.不写字段或类似的东西.
谢谢.
有一个在任务中执行的过程.我不希望其中一个同时执行.
这是检查任务是否已在运行的正确方法吗?
private Task task;
public void StartTask()
{
if (task != null && (task.Status == TaskStatus.Running || task.Status == TaskStatus.WaitingToRun || task.Status == TaskStatus.WaitingForActivation))
{
Logger.Log("Task has attempted to start while already running");
}
else
{
Logger.Log("Task has began");
task = Task.Factory.StartNew(() =>
{
// Stuff
});
}
}
Run Code Online (Sandbox Code Playgroud) 所以我只是学习新的Java 8,特别是lambdas和日期和时间api.我把它与scala进行比较.我的基本想法是找到命令行,流和并行流之间的执行时间差异.所以我决定创建一个Library应用程序并执行搜索,过滤,排序等操作.我创建了一个Library类,其中包含一个名为books的列表字段,并填充了1000本书.然后为搜索创建了一个功能界面,并在所有三种样式中进行了一些操作.一切都很好.我的代码是:
// Functional Interface
interface Search<T> {
public void search(T t);
}
// Library class
final Library library = new Library();
// This just creates some random book objects.
final List<Book> books = collectBooks();
final Search<List<Book>> parallelSearch = (bks) -> library.findAndPrintBooksParallel(bks);
// Parallel Operations
private void findAndPrintBooksParallel(List<Book> books) {
books.parallelStream()
.filter(b -> b.getAuthor().equals("J.K. Rowling"))
.sorted((x,y) -> x.getAuthor().compareTo(y.getAuthor()))
.map(Book::getIsbn)
.forEach(Library::waitAndPrintRecord);
}
Run Code Online (Sandbox Code Playgroud)
现在我尝试在scala中重新创建相同的程序,看看执行是否更快?令人惊讶的是scala不允许我进行并行排序(或者可能是我在这里无知).我的scala库是
// Again some random book objects as a list
val books = collectBooks
// Parallel operation
books.par filter(_.author …Run Code Online (Sandbox Code Playgroud) 为什么要forEach以随机顺序打印数字,同时collect始终按原始顺序收集元素,即使是从并行流中收集?
Integer[] intArray = {1, 2, 3, 4, 5, 6, 7, 8};
List<Integer> listOfIntegers = new ArrayList<>(Arrays.asList(intArray));
System.out.println("Parallel Stream: ");
listOfIntegers
.stream()
.parallel()
.forEach(e -> System.out.print(e + " "));
System.out.println();
// Collectors
List<Integer> l = listOfIntegers
.stream()
.parallel()
.collect(Collectors.toList());
System.out.println(l);
Run Code Online (Sandbox Code Playgroud)
输出:
Parallel Stream:
8 1 6 2 7 4 5 3
[1, 2, 3, 4, 5, 6, 7, 8]
Run Code Online (Sandbox Code Playgroud) java ×4
c# ×3
java-stream ×3
java-8 ×2
r ×2
.net ×1
asynchronous ×1
forkjoinpool ×1
lambda ×1
list ×1
php ×1
scala ×1
task-queue ×1
threadpool ×1