我有以下代码可以读取一个大文件,比如超过一百万行。我正在使用 Parallel 和 Linq 方法。有没有更好的方法来做到这一点?如果是,那么如何?
private static void ReadFile()
{
float floatTester = 0;
List<float[]> result = File.ReadLines(@"largedata.csv")
.Where(l => !string.IsNullOrWhiteSpace(l))
.Select(l => new { Line = l, Fields = l.Split(new[] { ',' }, StringSplitOptions.RemoveEmptyEntries) })
.Select(x => x.Fields
.Where(f => Single.TryParse(f, out floatTester))
.Select(f => floatTester).ToArray())
.ToList();
// now get your totals
int numberOfLinesWithData = result.Count;
int numberOfAllFloats = result.Sum(fa => fa.Length);
MessageBox.Show(numberOfAllFloats.ToString());
}
private static readonly char[] Separators = { ',', ' ' };
private static void ProcessFile() …Run Code Online (Sandbox Code Playgroud) 我意识到这在很大程度上取决于相关流程,但是否有经验法则?
假设我有一个名为的多线程程序progX,它提供一个命令行开关 ( --cpu) 来控制它可以使用的 CPU 数量。启动 40 个并行实例每个使用一个 CPU ( progX --cpu 1) 还是启动单个实例并告诉它使用 40 个 CPU ( progX --cpu 40)是否更快?
我正在使用在 python 进程之间multiprocessing.Queue传递 numpy 数组float64。这工作正常,但我担心它可能没有达到应有的效率。
根据 的文档multiprocessing,放置在 上的对象Queue将被腌制。调用picklenumpy 数组会产生数据的文本表示,因此空字节被 string 替换"\\x00"。
>>> pickle.dumps(numpy.zeros(10))
"cnumpy.core.multiarray\n_reconstruct\np0\n(cnumpy\nndarray\np1\n(I0\ntp2\nS'b'\np3\ntp4\nRp5\n(I1\n(I10\ntp6\ncnumpy\ndtype\np7\n(S'f8'\np8\nI0\nI1\ntp9\nRp10\n(I3\nS'<'\np11\nNNNI-1\nI-1\nI0\ntp12\nbI00\nS'\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00\\x00'\np13\ntp14\nb."
我担心这意味着我的数组被昂贵地转换为原始大小的 4 倍,然后在另一个过程中转换回。
有没有办法以原始未更改的形式通过队列传递数据?
我知道共享内存,但如果这是正确的解决方案,我不确定如何在其上构建队列。
谢谢!
目标
使用GNU Parallel将大的.gz文件拆分为子文件.由于服务器有16个CPU,因此创建16个子节点.每个孩子最多应包含N行.这里,N = 104,214,420行.儿童应该是.gz格式.
输入文件
硬件
码
版本1
zcat "${input_file}" | parallel --pipe -N 104214420 --joblog split_log.txt --resume-failed "gzip > ${input_file}_child_{#}.gz"
Run Code Online (Sandbox Code Playgroud)
三天后,工作还没完成.split_log.txt为空.输出目录中没有可见的子项.日志文件表明Parallel --block-size已从1 MB(默认值)增加到2 GB以上.这激发了我将代码更改为版本2.
版本2
# --block-size 3000000000 means a single record could be 3 GB long. Parallel will increase this value if needed.
zcat "${input_file}" | "${parallel}" --pipe -N 104214420 --block-size 3000000000 --joblog split_log.txt --resume-failed "gzip > ${input_file}_child_{#}.gz"
Run Code Online (Sandbox Code Playgroud)
这项工作已经运行了大约2个小时.split_log.txt为空.尚未在输出目录中看到子项.到目前为止,日志文件显示以下警告:
parallel: Warning: --blocksize >= 2G causes problems. Using 2G-1. …Run Code Online (Sandbox Code Playgroud) 我是张量流的初学者.我目前正在研究一个拥有2个GPU的系统,每个12GB.我想在两个GPU上实现模型并行性来训练大型模型.我一直在浏览整个互联网,SO,tensorflow文档等,我能够找到模型并行性及其结果的解释,但我没有找到一个小教程或小代码片段如何使用tensorflow实现它.我的意思是我们必须在每一层之后交换激活权利吗?那我们该怎么做呢?在tensorflow中是否有一种特定的或更简洁的方法来实现模型并行性?如果你可以建议我学习实现它的地方,或者使用'MODEL PARALLELISM'在多GPU上进行mnist训练这样的简单代码,将会非常有用.
注意:我已经完成了像CIFAR10中的数据并行 - 多gpu教程,但我没有找到任何模型并行的实现.
我有一个小代码片段如下:
import requests
import multiprocessing
header = {
'X-Location': 'UNKNOWN',
'X-AppVersion': '2.20.0',
'X-UniqueId': '2397123',
'X-User-Locale': 'en',
'X-Platform': 'Android',
'X-AppId': 'com.my_app',
'Accept-Language': 'en-ID',
'X-PushTokenType': 'GCM',
'X-DeviceToken': 'some_device_token'
}
BASE_URI = 'https://my_server.com/v2/customers/login'
def internet_resource_getter(post_data):
stuff_got = []
response = requests.post(BASE_URI, headers=header, json=post_data)
stuff_got.append(response.json())
return stuff_got
tokens = [{"my_token":'EAAOZAe8Q2rKYBAu0XETMiCZC0EYAddz4Muk6Luh300PGwGAMh26Bpw3AA6srcxbPWSTATpTLmvhzkUHuercNlZC1vDfL9Kmw3pyoQfpyP2t7NzPAOMCbmCAH6ftXe4bDc4dXgjizqnudfM0D346rrEQot5H0esW3RHGf8ZBRVfTtX8yR0NppfU5LfzNPqlAem9M5ZC8lbFlzKpZAZBOxsaz'},{"my_token":'EAAOZAe8Q2rKYBAKQetLqFwoTM2maZBOMUZA2w5mLmYQi1GpKFGZAxZCaRjv09IfAxxK1amZBE3ab25KzL4Bo9xvubiTkRriGhuivinYBkZAwQpnMZC99CR2FOqbNMmZBvLjZBW7xv6BwSTu3sledpLSGQvPIZBKmTv3930dBH8lazZCs3q0Q5i9CZC8mf8kYeamV9DED1nsg5PQZDZD'}]
pool = multiprocessing.Pool(processes=3)
pool_outputs = pool.map(internet_resource_getter, tokens)
pool.close()
pool.join()
Run Code Online (Sandbox Code Playgroud)
我所要做的就是将并行POST请求发送到终点,而每个POST都有一个不同的令牌,因为它的帖子正文.
parallel-processing multiprocessing python-2.7 python-requests grequests
我是java的新手,我想使用执行器服务或使用java中的任何其他方法并行化嵌套for循环.我想创建一些固定数量的线程,以便线程不会完全获取CPU.
for(SellerNames sellerNames : sellerDataList) {
for(String selleName : sellerNames) {
//getSellerAddress(sellerName)
//parallize this task
}
}
Run Code Online (Sandbox Code Playgroud)
sellerDataList = 1000的大小和sellerNames = 5000的大小.
现在我想创建10个线程并将相同的任务块分配给每个线程.这是为了我的sellerDataList,第一个线程应该获得500个名称的地址,第二个线程应该获得下一个500个名称的地址,依此类推.
做这份工作的最佳方法是什么?
我正在学习Python中的线程库。我不明白,如何并行运行两个线程?
这是我的python程序:
没有线程的程序(fibsimple.py)
def fib(n):
if n < 2:
return n
else:
return fib(n-1) + fib(n-2)
fib(35)
fib(35)
print "Done"
Run Code Online (Sandbox Code Playgroud)
运行时间:
$ time python fibsimple.py
Done
real 0m7.935s
user 0m7.922s
sys 0m0.008s
Run Code Online (Sandbox Code Playgroud)
带有线程的相同程序(fibthread.py)
from threading import Thread
def fib(n):
if n < 2:
return n
else:
return fib(n-1) + fib(n-2)
t1 = Thread(target = fib, args = (35, ))
t1.start()
t2 = Thread(target = fib, args = (35, ))
t2.start()
t1.join()
t2.join()
print "Done"
Run Code Online (Sandbox Code Playgroud)
运行时间:
$ …Run Code Online (Sandbox Code Playgroud) python parallel-processing python-multithreading python-multiprocessing
我正在使用Stream并行处理,并了解如果我使用平面阵列流,它会得到非常快速的处理.但如果我使用ArrayList,那么处理速度会慢一些.但是,如果我使用LinkedList或使用一些二进制树,处理速度会更慢.
所有听起来更像是流的可分割性,处理速度越快.这意味着阵列和数组列表在并行流的情况下最有效.这是真的吗?如果是这样,ArrayList如果我们想并行处理流,我们总是使用或者Array吗?如果是这样,如何使用LinkedList和BlockingQueue并行流?
另一件事是选择的中间函数的状态.如果我执行像无状态操作filter(),map(),性能高,但如果执行像国家提供充分的操作distinct(),sorted(),limit(),skip(),它需要大量的时间.再次,并行流变慢.这是否意味着我们不应该在并行流中使用状态全中间函数?如果是这样,那么解决这个问题的方法是什么?
我试图在一个非常大的数据集上运行一些东西.基本上,我想遍历文件夹中的所有文件并在其上运行fromJSON函数.但是,我希望它跳过产生错误的文件.我已经使用tryCatch构建了一个函数,但只有在我使用函数lappy而不是parLapply时才有效.
这是我的异常处理函数的代码:
readJson <- function (file) {
require(jsonlite)
dat <- tryCatch(
{
fromJSON(file, flatten=TRUE)
},
error = function(cond) {
message(cond)
return(NA)
},
warning = function(cond) {
message(cond)
return(NULL)
}
)
return(dat)
}
Run Code Online (Sandbox Code Playgroud)
然后我在包含JSON文件的完整路径的字符向量文件上调用parLapply :
dat<- parLapply(cl,files,readJson)
Run Code Online (Sandbox Code Playgroud)
当它到达一个未正确结束的文件时会产生错误,并且不会通过跳过有问题的文件来创建列表'dat'.这是readJson函数应该缓解的内容.
当我使用常规lapply,但它工作得很好.它会生成错误,但是,它仍然会跳过错误的文件来创建列表.
关于如何使用parLappy并行处理异常处理的任何想法,以便它会跳过有问题的文件并生成列表?
python ×2
bash ×1
c# ×1
distributed ×1
gnu-parallel ×1
grequests ×1
java ×1
java-8 ×1
java-stream ×1
linq ×1
numpy ×1
performance ×1
pickle ×1
python-2.7 ×1
r ×1
tensorflow ×1
threadpool ×1