- android - RelativeLayout 背景可绘制重叠内容
- android - 如何链接 cpufeatures lib 以获取 native android 库?
- java - OnItemClickListener 不起作用,但 OnLongItemClickListener 在自定义 ListView 中起作用
- java - Android 文件转字符串
我正在尝试使用 TPL Dataflow
实现数据处理管道.但是,我对数据流比较陌生,并不完全确定如何正确使用它来解决我要解决的问题。
问题:
我正在尝试遍历文件列表并处理每个文件以读取一些数据,然后进一步处理该数据。每个文件大概是700MB
至 1GB
在尺寸方面。每个文件包含 JSON
数据。为了并行处理这些文件而不是运行内存,我正在尝试使用 IEnumerable<>
与 yield return
然后进一步处理数据。
获得文件列表后,我想同时处理最多 4-5 个文件。我的困惑来自:
IEnumerable<>
和 yeild return
与 async/await
和数据流。偶遇this answer通过 svick , 但仍然不确定如何转换 IEnumerable<>
至 ISourceBlock
然后将所有 block 链接在一起并跟踪完成情况。producer
会非常快(通过文件列表),但是 consumer
将非常慢(处理每个文件 - 读取数据,反序列化 JSON
)。在这种情况下,如何跟踪完成情况。LinkTo
吗?数据 block 连接各种 block 的功能?或使用 OutputAvailableAsync()
等方法和 ReceiveAsync()
将数据从一个 block 传播到另一个 block 。代码:
private const int ProcessingSize= 4;
private BufferBlock<string> _fileBufferBlock;
private ActionBlock<string> _processingBlock;
private BufferBlock<DataType> _messageBufferBlock;
public Task ProduceAsync()
{
PrepareDataflow(token);
var bufferTask = ListFilesAsync(_fileBufferBlock, token);
var tasks = new List<Task> { bufferTask, _processingBlock.Completion };
return Task.WhenAll(tasks);
}
private async Task ListFilesAsync(ITargetBlock<string> targetBlock, CancellationToken token)
{
...
// Get list of file Uris
...
foreach(var fileNameUri in fileNameUris)
await targetBlock.SendAsync(fileNameUri, token);
targetBlock.Complete();
}
private async Task ProcessFileAsync(string fileNameUri, CancellationToken token)
{
var httpClient = new HttpClient();
try
{
using (var stream = await httpClient.GetStreamAsync(fileNameUri))
using (var sr = new StreamReader(stream))
using (var jsonTextReader = new JsonTextReader(sr))
{
while (jsonTextReader.Read())
{
if (jsonTextReader.TokenType == JsonToken.StartObject)
{
try
{
var data = _jsonSerializer.Deserialize<DataType>(jsonTextReader)
await _messageBufferBlock.SendAsync(data, token);
}
catch (Exception ex)
{
_logger.Error(ex, $"JSON deserialization failed - {fileNameUri}");
}
}
}
}
}
catch(Exception ex)
{
// Should throw?
// Or if converted to block then report using Fault() method?
}
finally
{
httpClient.Dispose();
buffer.Complete();
}
}
private void PrepareDataflow(CancellationToken token)
{
_fileBufferBlock = new BufferBlock<string>(new DataflowBlockOptions
{
CancellationToken = token
});
var actionExecuteOptions = new ExecutionDataflowBlockOptions
{
CancellationToken = token,
BoundedCapacity = ProcessingSize,
MaxMessagesPerTask = 1,
MaxDegreeOfParallelism = ProcessingSize
};
_processingBlock = new ActionBlock<string>(async fileName =>
{
try
{
await ProcessFileAsync(fileName, token);
}
catch (Exception ex)
{
_logger.Fatal(ex, $"Failed to process fiel: {fileName}, Error: {ex.Message}");
// Should fault the block?
}
}, actionExecuteOptions);
_fileBufferBlock.LinkTo(_processingBlock, new DataflowLinkOptions { PropagateCompletion = true });
_messageBufferBlock = new BufferBlock<DataType>(new ExecutionDataflowBlockOptions
{
CancellationToken = token,
BoundedCapacity = 50000
});
_messageBufferBlock.LinkTo(DataflowBlock.NullTarget<DataType>());
}
在上面的代码中,我没有使用 IEnumerable<DataType>
和 yield return
因为我不能将它与 async/await
一起使用.所以我将输入缓冲区链接到 ActionBlock<DataType>
这又发布到另一个队列。但是通过使用 ActionBlock<>
,我无法将其链接到下一个 block 进行处理,必须手动 Post/SendAsync
来自 ActionBlock<>
至 BufferBlock<>
.此外,在这种情况下,不确定如何跟踪完成情况。
这段代码有效,但是,我相信会有比这更好的解决方案,我可以链接所有 block (而不是 ActionBlock<DataType>
,然后从它发送消息到 BufferBlock<DataType>
)
另一种选择是转换 IEnumerable<>
至 IObservable<>
使用 Rx
, 但我又不太熟悉 Rx
并且不知道如何混合 TPL Dataflow
和 Rx
最佳答案
问题一
你插入一个IEnumerable<T>
使用 Post
将生产者添加到您的 TPL 数据流链中或 SendAsync
直接在消费者 block 上,如下:
foreach (string fileNameUri in fileNameUris)
{
await _processingBlock.SendAsync(fileNameUri).ConfigureAwait(false);
}
您还可以使用 BufferBlock<TInput>
,但在您的情况下,它实际上似乎是不必要的(甚至有害 - 请参阅下一部分)。
问题二
你更喜欢什么时候SendAsync
而不是 Post
?如果您的生产者运行速度快于 URI 的处理速度(并且您已表明是这种情况),并且您选择提供您的 _processingBlock
一个BoundedCapacity
,然后当 block 的内部缓冲区达到指定容量时,您的 SendAsync
将“挂起”直到缓冲区插槽释放,并且您的 foreach
循环将被限制。这种反馈机制会产生背压并确保您不会耗尽内存。
问题三
你绝对应该使用 LinkTo
在大多数情况下链接您的 block 的方法。不幸的是,由于 IDisposable
的相互作用,您的情况属于极端情况。和非常大的(潜在的)序列。所以你的完成将在缓冲区和处理 block 之间自动流动(由于 LinkTo
),但在那之后 - 你需要手动传播它。这很棘手,但可行。
我将用一个“Hello World”示例来说明这一点,其中生产者迭代每个字符,而消费者(这真的很慢)将每个字符输出到调试窗口。
备注:LinkTo
不存在。
// REALLY slow consumer.
var consumer = new ActionBlock<char>(async c =>
{
await Task.Delay(100);
Debug.Print(c.ToString());
}, new ExecutionDataflowBlockOptions { BoundedCapacity = 1 });
var producer = new ActionBlock<string>(async s =>
{
foreach (char c in s)
{
await consumer.SendAsync(c);
Debug.Print($"Yielded {c}");
}
});
try
{
producer.Post("Hello world");
producer.Complete();
await producer.Completion;
}
finally
{
consumer.Complete();
}
// Observe combined producer and consumer completion/exceptions/cancellation.
await Task.WhenAll(producer.Completion, consumer.Completion);
这个输出:
Yielded HHYielded eeYielded llYielded llYielded ooYielded Yielded wwYielded ooYielded rrYielded llYielded dd
As you can see from the output above, the producer is throttled and the handover buffer between the blocks never grows too large.
EDIT
You might find it cleaner to propagate completion via
producer.Completion.ContinueWith(
_ => consumer.Complete(), TaskContinuationOptions.ExecuteSynchronously
);
... 就在 producer
之后定义。这允许您稍微减少生产者/消费者耦合 - 但最后您仍然必须记住观察 Task.WhenAll(producer.Completion, consumer.Completion)
.
关于c# - 在 TPL 数据流中使用 async/await 和 yield return,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/35371931/
我需要将文本放在 中在一个 Div 中,在另一个 Div 中,在另一个 Div 中。所以这是它的样子: #document Change PIN
奇怪的事情发生了。 我有一个基本的 html 代码。 html,头部, body 。(因为我收到了一些反对票,这里是完整的代码) 这是我的CSS: html { backgroun
我正在尝试将 Assets 中的一组图像加载到 UICollectionview 中存在的 ImageView 中,但每当我运行应用程序时它都会显示错误。而且也没有显示图像。 我在ViewDidLoa
我需要根据带参数的 perl 脚本的输出更改一些环境变量。在 tcsh 中,我可以使用别名命令来评估 perl 脚本的输出。 tcsh: alias setsdk 'eval `/localhome/
我使用 Windows 身份验证创建了一个新的 Blazor(服务器端)应用程序,并使用 IIS Express 运行它。它将显示一条消息“Hello Domain\User!”来自右上方的以下 Ra
这是我的方法 void login(Event event);我想知道 Kotlin 中应该如何 最佳答案 在 Kotlin 中通配符运算符是 * 。它指示编译器它是未知的,但一旦知道,就不会有其他类
看下面的代码 for story in book if story.title.length < 140 - var story
我正在尝试用 C 语言学习字符串处理。我写了一个程序,它存储了一些音乐轨道,并帮助用户检查他/她想到的歌曲是否存在于存储的轨道中。这是通过要求用户输入一串字符来完成的。然后程序使用 strstr()
我正在学习 sscanf 并遇到如下格式字符串: sscanf("%[^:]:%[^*=]%*[*=]%n",a,b,&c); 我理解 %[^:] 部分意味着扫描直到遇到 ':' 并将其分配给 a。:
def char_check(x,y): if (str(x) in y or x.find(y) > -1) or (str(y) in x or y.find(x) > -1):
我有一种情况,我想将文本文件中的现有行包含到一个新 block 中。 line 1 line 2 line in block line 3 line 4 应该变成 line 1 line 2 line
我有一个新项目,我正在尝试设置 Django 调试工具栏。首先,我尝试了快速设置,它只涉及将 'debug_toolbar' 添加到我的已安装应用程序列表中。有了这个,当我转到我的根 URL 时,调试
在 Matlab 中,如果我有一个函数 f,例如签名是 f(a,b,c),我可以创建一个只有一个变量 b 的函数,它将使用固定的 a=a1 和 c=c1 调用 f: g = @(b) f(a1, b,
我不明白为什么 ForEach 中的元素之间有多余的垂直间距在 VStack 里面在 ScrollView 里面使用 GeometryReader 时渲染自定义水平分隔线。 Scrol
我想知道,是否有关于何时使用 session 和 cookie 的指南或最佳实践? 什么应该和什么不应该存储在其中?谢谢! 最佳答案 这些文档很好地了解了 session cookie 的安全问题以及
我在 scipy/numpy 中有一个 Nx3 矩阵,我想用它制作一个 3 维条形图,其中 X 轴和 Y 轴由矩阵的第一列和第二列的值、高度确定每个条形的 是矩阵中的第三列,条形的数量由 N 确定。
假设我用两种不同的方式初始化信号量 sem_init(&randomsem,0,1) sem_init(&randomsem,0,0) 现在, sem_wait(&randomsem) 在这两种情况下
我怀疑该值如何存储在“WORD”中,因为 PStr 包含实际输出。? 既然Pstr中存储的是小写到大写的字母,那么在printf中如何将其给出为“WORD”。有人可以吗?解释一下? #include
我有一个 3x3 数组: var my_array = [[0,1,2], [3,4,5], [6,7,8]]; 并想获得它的第一个 2
我意识到您可以使用如下方式轻松检查焦点: var hasFocus = true; $(window).blur(function(){ hasFocus = false; }); $(win
我是一名优秀的程序员,十分优秀!