标签: bigdata

hadoop map减少二次排序

谁能解释我在hadoop中如何进行二次排序?
为什么必须使用GroupingComparator它以及它在hadoop中如何工作?

我正在浏览下面给出的链接,并怀疑groupcomapator如何工作.
任何人都可以解释一下分组比较器的工作原理吗?

http://www.bigdataspeak.com/2013/02/hadoop-how-to-do-secondary-sort-on_25.html

hadoop mapreduce bigdata hadoop-partitioning

21
推荐指数
3
解决办法
3万
查看次数

Spark RDD - 它们是如何工作的

我有一个小型Scala程序,可以在单个节点上运行.但是,我正在扩展它,因此它在多个节点上运行.这是我的第一次尝试.我只是想了解RDD如何在Spark中工作,所以这个问题是基于理论的,可能不是100%正确.

假设我创建了一个RDD: val rdd = sc.textFile(file)

现在,一旦我这样做了,这是否意味着文件at file现在在节点之间进行分区(假设所有节点都可以访问文件路径)?

其次,我想计算RDD中的对象数量(足够简单),但是,我需要在需要应用于RDD中的对象的计算中使用该数字 - 伪代码示例:

rdd.map(x => x / rdd.size)
Run Code Online (Sandbox Code Playgroud)

假设有100个对象rdd,并且说有10个节点,因此每个节点有10个对象的计数(假设这是RDD概念的工作方式),现在当我调用该方法时,每个节点将使用rdd.sizeas 执行计算10还是100?因为,总的来说,RDD是大小100但在每个节点上本地只是10.我是否需要在进行计算之前制作广播变量?这个问题与下面的问题有关.

最后,如果我转换到RDD,例如rdd.map(_.split("-")),然后我想要新size的RDD,我是否需要在RDD上执行操作,例如count(),所以所有信息都被发送回驱动程序节点?

scala distributed-computing bigdata apache-spark rdd

21
推荐指数
2
解决办法
4777
查看次数

在巨大的文字上找到最重复的短语

我有大量的文本数据.我的整个数据库都是UTF-8的文本格式

我需要在我的整个文本数据上列出最重复的短语.

例如,我的愿望输出如下:

{
  'a': 423412341,
  'this': 423412341,
  'is': 322472341,
  'this is': 222472341,
  'this is a': 122472341,
  'this is a my': 5235634
}
Run Code Online (Sandbox Code Playgroud)

处理和存储每个短语占用巨大的数据库.例如存储在MySQL或MongoDB中.问题是有没有更有效的数据库或算法来找到这个结果?Solr,Elasticsearch等......

我想我每个短语最多10个单词对我有好处.

search text full-text-search bigdata

21
推荐指数
1
解决办法
2008
查看次数

运行glmnet()的大矩阵

我有一个问题,用宽数据集运行glmnet套索.我的数据N = 50,但p> 49000,所有因素.所以要运行glmnet,我必须创建一个model.matrix,但是当我调用model.matrix(formula,data)时,我的内存耗尽,其中formula = Class~.

作为一个工作示例,我将生成一个数据集:

data <- matrix(rep(0,50*49000), nrow=50)
for(i in 1:50) {
x = rep(letters[2:8], 7000)
y = sample(x=1:49000, size=49000)
data[i,] <- x[y]
}

data <- as.data.frame(data)
x = c(rep('A', 20), rep('B', 15), rep('C', 15))
y = sample(x=1:50, size=50)
class = x[y]
data <- cbind(data, class)
Run Code Online (Sandbox Code Playgroud)

之后,我尝试创建一个model.matrix进入glmnet.

  formula <- as.formula(class ~ .)
  X = model.matrix(formula, data)
  model <- cv.glmnet(X, class, standardize=FALSE, family='multinomial', alpha=1, nfolds=10)
Run Code Online (Sandbox Code Playgroud)

在最后一步(X = model.matrix ...),我的内存不足.我能做什么?

r lasso-regression bigdata glmnet model.matrix

20
推荐指数
2
解决办法
8773
查看次数

读取固定宽度的大数据

如何读取固定宽度格式的大数据?我读了这个问题并尝试了一些提示,但所有答案都是针对分隔数据(如.csv),而这不是我的理由.数据有558MB,我不知道有多少行.

我正在使用:

dados <- read.fwf('TS_MATRICULA_RS.txt', width=c(5, 13, 14, 3, 3, 5, 4, 6, 6, 6, 1, 1, 1, 4, 3, 2, 9, 3, 2, 9, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1,
    1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 3, 4, 11, 9, 2, 3, 9, 3, 2, 9, 9, 1, 1, 1, 1, 2, 1, 1, 1, 1, 1, 1, 1, 1), …
Run Code Online (Sandbox Code Playgroud)

r bigdata

20
推荐指数
3
解决办法
8998
查看次数

将大量数据从Cassandra导出到CSV

我正在使用Cassandra 2.0.9存储相当大的数据,比如100Gb,在一个列族中.我想快速将此数据导出为CSV.我试过了:

  • sstable2json - 它生成相当大的json文件,难以解析 - 因为工具将数据放在一行并使用复杂的模式(例如300Mb数据文件=〜2Gb json),转储需要花费大量时间,而Cassandra喜欢改变源文件名根据其内部机制
  • COPY - 在相当快的EC2实例上导致大量记录的超时
  • 捕获 - 如上所述,导致超时
  • 用分页读取 - 我使用了timeuuid,但它每秒返回大约1,5k条记录

我使用Amazon Ec2实例,具有快速存储,15 Gb RAM和4个内核

对于从Cassandra到CSV的数据导出千兆字节有什么更好的选择吗?

csv bigdata cassandra cassandra-2.0

20
推荐指数
1
解决办法
1万
查看次数

pyspark mapPartitions函数如何工作?

所以我试图用Python(Pyspark)学习Spark.我想知道这个功能是如何mapPartitions工作的.这就是输入它所带来的输出和输出.我在互联网上找不到任何合适的例子.可以说,我有一个包含列表的RDD对象,如下所示.

[ [1, 2, 3], [3, 2, 4], [5, 2, 7] ] 
Run Code Online (Sandbox Code Playgroud)

我想从所有列表中删除元素2,我将如何使用它mapPartitions.

python scala bigdata apache-spark

20
推荐指数
2
解决办法
3万
查看次数

根据worker,core和DataFrame大小确定Spark分区的最佳数量

Spark-land中有几个相似但又不同的概念,围绕着如何将工作分配到不同的节点并同时执行.具体来说,有:

  • Spark Driver节点(sparkDriverCount)
  • Spark群集可用的工作节点数(numWorkerNodes)
  • Spark执行器的数量(numExecutors)
  • 由所有工人/执行者同时操作的DataFrame(dataFrame)
  • dataFrame(numDFRows)中的行数
  • dataFrame(numPartitions)上的分区数
  • 最后,每个工作节点上可用的CPU核心数量(numCpuCoresPerWorker)

相信所有Spark集群都有一个且只有一个 Spark Driver,然后是0+个工作节点.如果我错了,请先纠正我!假设我或多或少是正确的,让我们在这里锁定一些变量.假设我们有一个带有1个驱动程序和4个工作节点的Spark集群,每个工作节点上有4个CPU核心(因此总共有16个CPU核心).所以这里的"给定"是:

sparkDriverCount = 1
numWorkerNodes = 4
numCpuCores = numWorkerNodes * numCpuCoresPerWorker = 4 * 4 = 16
Run Code Online (Sandbox Code Playgroud)

鉴于作为设置,我想知道如何确定一些事情.特别:

  • numWorkerNodes和之间有什么关系numExecutors?是否有一些已知/普遍接受的工人与遗嘱执行人的比例?有没有办法确定numExecutors给定numWorkerNodes(或任何其他输入)?
  • 是否已知/普遍接受/最佳比率numDFRowsnumPartitions?如何根据dataFrame?的大小计算"最佳"分区数?
  • 我从其他工程师那里得知,一般的"经验法则"是:numPartitions = numWorkerNodes * numCpuCoresPerWorker那有什么道理吗?换句话说,它规定每个CPU核心应该有一个分区.

partitioning distributed-computing bigdata apache-spark spark-dataframe

20
推荐指数
1
解决办法
1万
查看次数

单件读取CSV文件的策略?

我在计算机上有一个中等大小的文件(4GB CSV),没有足够的RAM来读取它(在64位Windows上为8GB).在过去,我只是将它加载到一个集群节点上并将其读入,但我的新集群似乎任意将进程限制为4GB的RAM(尽管硬件每台机器有16GB),所以我需要一个短期修复.

有没有办法将CSV文件的一部分读入R以适应可用的内存限制?这样我一次可以读取文件的三分之一,将其子集化为我需要的行和列,然后在下一个三分之一读取?

感谢评论者指出我可以使用一些大内存技巧读取整个文件: 快速读取非常大的表作为R中的数据帧

我可以想到其他一些解决方法(例如在一个好的文本编辑器中打开,删掉2/3的观察结果,然后加载R),但是如果可能的话我宁愿避免使用它们.

因此,阅读它看起来仍然是现在最好的方法.

r bigdata

19
推荐指数
2
解决办法
2万
查看次数

Haskell:我可以在同一个懒惰列表上执行多次折叠而不将列表保留在内存中吗?

我的背景是生物信息学,特别是新一代测序,但问题是通用的; 所以我将使用日志文件作为示例.

该文件非常大(千兆字节大,压缩,因此它不适合内存),但很容易解析(每行都是一个条目),所以我们可以轻松地写出如下内容:

parse :: Lazy.ByteString -> [LogEntry]
Run Code Online (Sandbox Code Playgroud)

现在,我有很多我想从日志文件中计算的统计信息.编写单独的函数是最简单的,例如:

totalEntries = length
nrBots = sum . map fromEnum . map isBotEntry
averageTimeOfDay = histogram . map extractHour
Run Code Online (Sandbox Code Playgroud)

所有这些都是这种形式foldl' k z . map f.

问题是,如果我尝试以最自然的方式使用它们,比如

main = do
    input <- Lazy.readFile "input.txt"
    let logEntries = parse input
        totalEntries' = totalEntries logEntries
        nrBots' = nrBots logEntries
        avgTOD = averageTimeOfDay logEntries
    print totalEntries'
    print nrBots'
    print avgTOD
Run Code Online (Sandbox Code Playgroud)

这会将整个列表分配到内存中,这不是我想要的.我希望折叠同步完成,以便可以对垃圾收集进行垃圾收集.如果我只计算一个统计量,就会发生这种情况.

我可以写一个这样做的大函数,但它是不可组合的代码.

或者,这就是我一直在做的,我分别运行每个传递,但每次重新加载和解压缩文件.

performance haskell lazy-evaluation bigdata

19
推荐指数
2
解决办法
809
查看次数