- html - 出于某种原因,IE8 对我的 Sass 文件中继承的 html5 CSS 不友好?
- JMeter 在响应断言中使用 span 标签的问题
- html - 在 :hover and :active? 上具有不同效果的 CSS 动画
- html - 相对于居中的 html 内容固定的 CSS 重复背景?
我有一个包含 500,000 个元素的列表和一个包含 20 个消费者的队列。消息以不同的速度处理(1、15、30、60 秒;3、50 分钟;3、16 小时或更长时间。24 小时为超时)。我需要消费者的响应才能对数据进行一些处理。我将为此使用 Scala Future
和基于事件的 onComplete
。
为了不淹没队列,我想先向队列发送 30 条消息:20 条将由消费者挑选,10 条将在队列中等待。当其中一个 Future
完成时,我想向队列发送另一条消息。你能告诉我如何实现这一目标吗?这可以用 Akka Streams 完成吗?
这是错误的,我只是想让你知道我想要什么:
private def sendMessage(ids: List[String]): Unit = {
val id = ids.head
val futureResult = Future {
//send id among some message to the queue
}.map { result =>
//process the response
}
futureResult.onComplete { _ =>
sendMessage(ids.tail)
}
}
def migrateAll(): Unit = {
val ids: List[String] = //get IDs from the DB
sendMessage(ids)
}
最佳答案
下面是一个使用 Akka Streams 对您的用例进行建模的简单示例。
让我们将处理定义为一个接受 String
并返回 Future[String]
的方法:
def process(id: String): Future[String] = ???
然后我们从包含 500,000 个 String
元素的 List
创建一个 Source
并使用 mapAsync
将元素提供给处理方法。并行级别设置为 20,这意味着在任何时间点运行的 Future
不超过 20 个。当每个 Future
完成时,我们执行额外的处理并打印结果:
Source((1 to 500000).map(_.toString).toList)
.mapAsync(parallelism = 20)(process)
// do something with the result of the Future; here we create a new string
// that begins with "Processed: "
.map(s => s"Processed: $s")
.runForeach(println)
您可以在 documentation 中阅读有关 mapAsync
的更多信息.
关于multithreading - 等待 Scala Future 完成并继续下一个,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/46713131/
从 Redis 获取消息时,onDone:(){print('done')} 从未起作用。 import 'package:dartis/dartis.dart' as redis show PubS
昨天我玩了一些vim脚本,并设法通过循环来对当前输入的内容进行状态栏预测(请参见屏幕截图(灰色+黄色栏))。 问题是,我不记得我是怎么得到的,也找不到我用于该vim魔术的代码片段(我记得它很简单):它
我尝试加载 bash_completion在我的 bash (3.2.25) 中,它不起作用。没有消息等。我在我的 .bashrc 中使用了以下内容 if [ -f ~/.bash_completio
我正在尝试构建一个 bash 完成例程,它将建议命令行标志和合适的标志值。例如在下面 fstcompose 命令我想比赛套路先建议 compose_filter= 标志,然后建议来自 [alt_seq
当我尝试在重定向符号后完成路径时,bash 完成的行为就好像它仍在尝试在重定向之前完成命令的参数一样。 例如: dpkg -l > /med标签 通过在 /med 之后点击 Tab我希望它完成通往 /
我的类中有几个 CAKeyframeAnimation 对象。 他们都以 self 为代表。 在我的animationDidStop函数中,我如何知道调用来自哪里? 是否有任何变量可以传递给 CAKe
我有一个带有 NSDateFormatter 的 NSTextField。格式化程序接受“mm/dd/yy”。 可以自动补全日期吗?因此,用户可以输入“mm”,格式化程序将完成当前月份和年份。 最佳答
有一个解决方案可以使用以下方法完成 NSTextField : - (NSArray *)control:(NSControl *)control textView:(NSTextView *)tex
我正在阅读 Passport 的文档,我注意到 serialize()和 deserialize() done()被调用而不被返回。 但是,当使用 passport.use() 设置新策略时在回调函数
在 ubuntu 11.10 上的 Firefox 8.0 中,尽管 img.complete 为 false,但仍会调用 onload 函数 draw。我设法用 setTimeout hack 解决
假设我有两个与两个并行执行的计算相对应的 future 。我如何等到第一个 future 准备好?理想情况下,我正在寻找类似于Python asyncio's wait且参数为return_when=
我正在寻找一种 Java 7 数据结构,其行为类似于 java.util.Queue,并且还具有“最终项目已被删除”的概念。 例如,应可以表达如下概念: while(!endingQueue.isFi
这是一个简单的问题。 if ($('.dataTablePageList')) { 我想做的是执行一个 if 语句,该语句表示如果具有 dataTablesPageList 类的对象也具有 menu
我用replaceWith批量替换了许多div中的html。替换后,我使用 jTruncate 来截断文本。然而它不起作用,因为在执行时,replaceWith 还没有完成。 我尝试了回调技巧 ( H
有没有办法调用 javascript 表单 submit() 函数或 JQuery $.submit() 函数并确保它完成提交过程?具体来说,在一个表单中,我试图在一个 IFrame 中提交一个表单。
我有以下方法: function animatePortfolio(fadeElement) { fadeElement.children('article').each(function(i
我刚刚开始使用 AndEngine, 我正在像这样移动 Sprite : if(pValueY < 0 && !jumping) { jumping =
我正在使用 asynctask 来执行冗长的操作,例如数据库读取。我想开始一个新 Activity 并在所有异步任务完成后呈现其内容。实现这一目标的最佳方法是什么? 我知道 onPostExecute
我有一个脚本需要命令名称和该命令的参数作为参数。 所以我想编写一个完成函数来完成命令的名称并完成该命令的参数。 所以我可以这样完成命令的名称 if [[ "$COMP_CWORD" == 1 ]];
我的应用程序有一个相当奇怪的行为。我在 BOOT_COMPLETE 之后启动我的应用程序,因此在我启动设备后它是可见的。 GUI 响应迅速,一切正常,直到我调用 finish(),按下按钮时,什么都没
我是一名优秀的程序员,十分优秀!