标签: bigdata

创建可比较且灵活的对象指纹

我的情况

说我有成千上万的对象,在这个例子中可能是电影.

我以很多不同的方式解析这些电影,收集有关每个电影的参数,关键字和统计数据.我们称他们为钥匙.我还为每个键分配一个权重,范围从0到1,具体取决于频率,相关性,强度,分数等.

作为一个例子,这里是电影世界末日的几个键和权重:

"Armageddon"
------------------
disaster       0.8
bruce willis   1.0
metascore      0.2
imdb score     0.4
asteroid       1.0
action         0.8
adventure      0.9
...            ...
Run Code Online (Sandbox Code Playgroud)

可能有成千上万的这些键和重量,为清楚起见,这是另一部电影:

"The Fast and the Furious"
------------------
disaster       0.1
bruce willis   0.0
metascore      0.5
imdb score     0.6
asteroid       0.0
action         0.9
adventure      0.6
...            ...
Run Code Online (Sandbox Code Playgroud)

我把它称为电影的指纹,我想用它们在我的数据库中查找类似的电影.

我还想象如果我愿意,可以插入除电影之外的其他内容,如文章或Facebook个人资料,并为其指定指纹.但那不应该影响我的问题.

我的问题

所以我已经走到了这一步,但现在我觉得这部分很棘手.我想把上面的指纹变成容易比较和快速的东西.我尝试创建一个数组,其中index 0= disaster,1= bruce willis,2= metascore,它们的值是权重.

我上面的两部电影就是这样的:

[ 0.8 , 1.0 , 0.2 , ... ] …
Run Code Online (Sandbox Code Playgroud)

c# sql algorithm data-mining bigdata

11
推荐指数
1
解决办法
1204
查看次数

将大型csv文件中的小型随机样本加载到R数据帧中

要处理的csv文件不适合内存.如何读取它的~20K随机线来对所选数据帧进行基本统计?

csv random r bigdata dataframe

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

大数据上的增量PCA

我刚刚尝试使用sklearn.decomposition中的IncrementalPCA,但它之前就像PCA和RandomizedPCA一样抛出了MemoryError.我的问题是,我试图加载的矩阵太大而无法放入RAM中.现在它作为shape~(1000000,1000)的数据集存储在hdf5数据库中,所以我有1.000.000.000 float32值.我认为IncrementalPCA批量加载数据,但显然它试图加载整个数据集,这没有帮助.这个库是如何使用的?hdf5格式是问题吗?

from sklearn.decomposition import IncrementalPCA
import h5py

db = h5py.File("db.h5","r")
data = db["data"]
IncrementalPCA(n_components=10, batch_size=1).fit(data)
Traceback (most recent call last):
  File "<stdin>", line 1, in <module>
  File "/software/anaconda/2.3.0/lib/python2.7/site-packages/sklearn/decomposition/incremental_pca.py", line 165, in fit
    X = check_array(X, dtype=np.float)
  File "/software/anaconda/2.3.0/lib/python2.7/site-packages/sklearn/utils/validation.py", line 337, in check_array
    array = np.atleast_2d(array)
  File "/software/anaconda/2.3.0/lib/python2.7/site-packages/numpy/core/shape_base.py", line 99, in atleast_2d
    ary = asanyarray(ary)
  File "/software/anaconda/2.3.0/lib/python2.7/site-packages/numpy/core/numeric.py", line 514, in asanyarray
    return array(a, dtype, copy=False, order=order, subok=True)
  File "h5py/_objects.pyx", line 54, in h5py._objects.with_phil.wrapper (-------src-dir-------/h5py/_objects.c:2458)
  File "h5py/_objects.pyx", line 55, in h5py._objects.with_phil.wrapper …
Run Code Online (Sandbox Code Playgroud)

python hdf5 bigdata pca scikit-learn

11
推荐指数
1
解决办法
6639
查看次数

大数据集成测试最佳实践

我正在寻找有关基于AWS的数据提取管道的最佳实践的一些资源,该管道使用Kafka,风暴,火花(流和批处理),使用各种微服务来读取和写入Hbase以暴露数据层.对于我的本地环境,我正在考虑创建docker或vagrant图像,这将允许我与env进行交互.我的问题就是如何为一个更接近生产的功能性端到端环境站起来的东西,这种下降方式就是拥有一个永远在线的环境,但这会变得昂贵.就perf性环境而言,似乎我可能不得不提出并拥有可以拥有"世界的运行"的服务帐户,但其他帐户将通过计算资源受到限制,因此它们不会压倒集群.

我很好奇其他人如何处理同样的问题,如果我正在考虑这个问题.

bigdata apache-spark apache-storm

11
推荐指数
1
解决办法
559
查看次数

如何跨分区平衡我的数据?

编辑:答案有帮助,但我在Spark中描述了我的解决方案:memoryOverhead问题.


我有一个带有202092分区的RDD,它读取其他人创建的数据集.我可以手动看到数据在分区之间没有平衡,例如它们中的一些有0个图像而其他有4k,而平均值是432.当处理数据时,我收到了这个错误:

Container killed by YARN for exceeding memory limits. 16.9 GB of 16 GB physical memory used. Consider boosting spark.yarn.executor.memoryOverhead.
Run Code Online (Sandbox Code Playgroud)

而memoryOverhead已经提升了.我觉得有些尖峰正在发生,这使得Yarn杀死我的容器,因为尖峰溢出了指定的边界.

那么我该怎么做才能确保我的数据在各个分区之间(大致)平衡?


我的想法是repartition()会工作,它会调用shuffling:

dataset = dataset.repartition(202092)
Run Code Online (Sandbox Code Playgroud)

但是我得到了同样的错误,尽管有编程指南的指示:

重新分区(numPartitions)

随机重新调整RDD中的数据以创建更多或更少的分区并在它们之间进行平衡.这总是随机播放网络上的所有数据.


检查我的玩具示例:

data = sc.parallelize([0,1,2], 3).mapPartitions(lambda x: range((x.next() + 1) * 1000))
d = data.glom().collect()
len(d[0])     # 1000
len(d[1])     # 2000
len(d[2])     # 3000
repartitioned_data = data.repartition(3)
re_d = repartitioned_data.glom().collect()
len(re_d[0])  # 1854
len(re_d[1])  # 1754 …
Run Code Online (Sandbox Code Playgroud)

python hadoop distributed-computing bigdata apache-spark

11
推荐指数
1
解决办法
5898
查看次数

Spark的KMeans无法处理bigdata吗?

KMeans有几个参数用于训练,初始化模式默认为kmeans ||.问题是它快速(少于10分钟)前进到前13个阶段,然后完全挂起,不会产生错误!

再现问题的最小示例(如果我使用1000点或随机初始化,它将成功):

from pyspark.context import SparkContext

from pyspark.mllib.clustering import KMeans
from pyspark.mllib.random import RandomRDDs


if __name__ == "__main__":
    sc = SparkContext(appName='kmeansMinimalExample')

    # same with 10000 points
    data = RandomRDDs.uniformVectorRDD(sc, 10000000, 64)
    C = KMeans.train(data, 8192,  maxIterations=10)    

    sc.stop()
Run Code Online (Sandbox Code Playgroud)

这项工作什么都不做(它没有成功,失败或进展......),如下所示."执行者"选项卡中没有活动/失败的任务.Stdout和Stderr Logs没有特别有趣的东西:

在此输入图像描述

如果我使用k=81,而不是8192,它将成功:

在此输入图像描述

请注意,这两个电话takeSample(),不应该是一个问题,因为有在随机初始化的情况下打了两次电话.

那么,发生了什么?Spark的Kmeans 无法扩展吗?有人知道吗?你可以重现吗?


如果这是一个内存问题,我会像以前一样得到警告和错误.

注意:placeybordeaux的注释基于在客户端模式下执行作业,其中驱动程序的配置无效,导致退出代码143等(请参阅编辑历史记录),而不是群集模式,其中根本没有报告错误,应用程序只是挂起.


从零到323:为什么Spark Mllib KMeans算法非常慢?是相关的,但我认为他目睹了一些进展,而我的确悬而未决,我确实发表评论......

在此输入图像描述

python bigdata k-means apache-spark apache-spark-mllib

11
推荐指数
1
解决办法
2645
查看次数

为什么Spark的OneHotEncoder默认删除最后一个类别?

我想了解Spark的OneHotEncoder默认情况下放弃最后一个类别的理性.

例如:

>>> fd = spark.createDataFrame( [(1.0, "a"), (1.5, "a"), (10.0, "b"), (3.2, "c")], ["x","c"])
>>> ss = StringIndexer(inputCol="c",outputCol="c_idx")
>>> ff = ss.fit(fd).transform(fd)
>>> ff.show()
+----+---+-----+
|   x|  c|c_idx|
+----+---+-----+
| 1.0|  a|  0.0|
| 1.5|  a|  0.0|
|10.0|  b|  1.0|
| 3.2|  c|  2.0|
+----+---+-----+
Run Code Online (Sandbox Code Playgroud)

默认情况下,OneHotEncoder将删除最后一个类别:

>>> oe = OneHotEncoder(inputCol="c_idx",outputCol="c_idx_vec")
>>> fe = oe.transform(ff)
>>> fe.show()
+----+---+-----+-------------+
|   x|  c|c_idx|    c_idx_vec|
+----+---+-----+-------------+
| 1.0|  a|  0.0|(2,[0],[1.0])|
| 1.5|  a|  0.0|(2,[0],[1.0])|
|10.0|  b|  1.0|(2,[1],[1.0])|
| 3.2|  c|  2.0|    (2,[],[])|
+----+---+-----+-------------+ …
Run Code Online (Sandbox Code Playgroud)

machine-learning bigdata apache-spark pyspark one-hot-encoding

11
推荐指数
1
解决办法
2673
查看次数

从关系数据库迁移到大数据

目前,我在Google云平台上托管了一个应用程序,该应用程序提供网络分析并提供会话活动(点击,下载等),并将该网络活动与网络注册联系起来.

目前,我们将所有点击和会话配置文件数据存储在MySQL中,并使用SQL查询生成聚合和每用户报告,但随着数据量的增长,我们看到查询响应真正减慢这反过来减慢了页面加载时间.

在调查我们可以解决这个问题的方法时,我们已经研究了Google云平台上可用的工具,如Dataproc和Dataflow以及NoSQL解决方案,但是,我很难理解如何将我们当前的解决方案应用于任何这些解决方案.

目前,我们对数据模式的概念如下:

User table
- id
- name
- email

Profile table (web browser/device)
- id
- user id
- user agent string

Session table
- id
- profile id
- session string

Action table
- id
- session id
- action type
- action details
- timestamp
Run Code Online (Sandbox Code Playgroud)

根据我的研究,我对什么是最佳解决方案的理解是将动作数据存储在NoT数据库解决方案中,如BigTable,它将数据提供给DataProc或DataFlow等生成报告的解决方案.但是,鉴于我们当前的架构是高度关系结构,似乎删除了转向NoSQL解决方案的选项,因为我的所有研究表明您不应该将关系数据移动到NoSQL解决方案.

我的问题是,我对如何正确应用这些工具的理解是什么?或者有更好的解决方案吗?是否有必要考虑远离MySQL?如果没有,有哪些解决方案可以让我们在后台预处理/生成报告数据?

mysql performance analytics bigdata google-cloud-platform

11
推荐指数
1
解决办法
920
查看次数

Elasticsearch部分批量更新

我在elasticsearch中有6kk的数据需要更新.我必须使用PHP.我在文档中搜索,我发现了这个,批量索引,但这不保留以前的数据.

我有:

[
  {
    'name': 'Jonatahn',
    'age' : 21
  }
]
Run Code Online (Sandbox Code Playgroud)

我要更新的代码:

$params =[
    "index" => "customer",
    "type" => "doc",
    "body" => [
        [
            "index" => [
                "_index" => "customer",
                "_type" => "doc",
                "_id" => "09310451939"
            ]
        ],
        [
            "name" => "Jonathan"
        ]
    ]
];

$client->bulk($params);
Run Code Online (Sandbox Code Playgroud)

当我发送时,['name' => 'Jonathan'] 我希望名称将更新并保持年龄,但年龄已被删除.当然,我仍然可以按数据更新数据,但这需要很长时间,还有另一种方法吗?

php json bulk bigdata elasticsearch

11
推荐指数
1
解决办法
2456
查看次数

Facebook等网站用什么格式存储个人资料的数据?

我最近开始处理存储在XML文件中的大量数据.我一直想知道Facebook和其他网络站点如何存储与个人配置文件相关的所有信息(名称,个人资料图片,墙上帖子等),我觉得XML绝对不是存储这么多信息的最佳方式.我已经尝试用谷歌查找有关它的信息,但没有太多的运气.

Facebook等大型网站如何存储和处理如此多的数据?我真的想读一下这个,所以如果你知道任何好的网站,请告诉我!

xml database storage facebook bigdata

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