标签: bigdata

如何从SQL查询创建大型pandas数据框而不会耗尽内存?

我无法从MS SQL Server数据库查询大于500万条记录的表.我希望能够选择所有记录,但在选择大量数据到内存时,我的代码似乎失败了.

这有效:

import pandas.io.sql as psql
sql = "SELECT TOP 1000000 * FROM MyTable" 
data = psql.read_frame(sql, cnxn)
Run Code Online (Sandbox Code Playgroud)

...但这不起作用:

sql = "SELECT TOP 2000000 * FROM MyTable" 
data = psql.read_frame(sql, cnxn)
Run Code Online (Sandbox Code Playgroud)

它返回此错误:

File "inference.pyx", line 931, in pandas.lib.to_object_array_tuples
(pandas\lib.c:42733) Memory Error
Run Code Online (Sandbox Code Playgroud)

我在这里读到,从csv文件创建数据帧时存在类似的问题,并且解决方法是使用'iterator'和'chunksize'参数,如下所示:

read_csv('exp4326.csv', iterator=True, chunksize=1000)
Run Code Online (Sandbox Code Playgroud)

是否有类似的SQL数据库查询解决方案?如果没有,首选的解决方法是什么?我是否需要通过其他方法读取块中的记录?我在这里阅读了一些关于在pandas中处理大型数据集的讨论,但执行SELECT*查询似乎需要做很多工作.当然有一种更简单的方法.

python sql bigdata pandas

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

机器学习和大数据

一开始我想描述一下我目前的立场和想要实现的目标.

我是一名处理机器学习的研究员.到目前为止,已经完成了几个涵盖机器学习算法和社交网络分析的理论课程,因此获得了一些理论概念,可用于实现机器学习算法并提供实际数据.

在简单的例子中,算法运行良好,运行时间可以接受,而大数据表示如果尝试在我的PC上运行算法,则会出现问题.关于软件我有足够的经验来实现文章中的任何算法或使用任何语言或IDE设计我自己(迄今为止使用过Matlab,Java与Eclipse,.NET ......)但到目前为止还没有多少设置经验 - 基础设施.我已经开始了解Hadoop,NoSQL数据库等,但我不确定哪种策略最好考虑学习时间限制.

最终目标是能够建立一个分析大数据的工作平台,重点是实现我自己的机器学习算法,并将所有这些算法集中到生产中,为处理大数据解决有用的问题做好准备.

由于主要关注的是实现机器学习算法,我想问一下是否有任何现有的运行平台,提供足够的CPU资源来提供大量数据,上传自己的算法并简单地处理数据而不考虑分布式处理.

然而这种平台的存在与否,我想获得足够的大,以便能够在一个团队中可能投产于特定的客户需求定制的整个系统工作的图片.例如,零售商希望分析每日购买,因此所有日常记录必须上传到某些基础设施,足以通过使用自定义机器学习算法处理数据.

将上述所有问题都纳入简单的问题:如何设计一个针对现实生活问题的自定义数据挖掘解决方案,主要关注机器学习算法并在可能的情况下通过使用现有基础设施投入生产,如果不是,则设计分布式系统(通过使用Hadoop或任何框架).

我非常感谢有关书籍或其他有用资源的任何建议或建议.

machine-learning bigdata

34
推荐指数
2
解决办法
7015
查看次数

我们可以使用什么方法来重塑非常大的数据集?

当由于非常大的数据计算将花费很长时间并且因此我们不希望它们崩溃时,事先知道要使用哪种重塑方法将很有价值。

最近,关于性能的数据重塑方法已得到进一步发展,例如data.table::dcasttidyr::spread。尤其dcast.data.table似乎设置了基调[1][2][3][4]。这使得基准R中的其他方法reshape显得过时且几乎无用[5]

理论

但是,我听说对于reshape大型数据集(可能是超出RAM的数据集)来说,这仍然是无与伦比的,因为它是唯一可以处理它们的方法,因此它仍然存在。与reshape2::dcast此相关的崩溃报告支持这一点 [6]。至少有一个参考文献给出了一个暗示,它reshape()可能确实比reshape2::dcast真正的“大杂烩” [7]具有优势。

方法

为此寻求证据,我认为值得花时间进行一些研究。所以我做了不同大小的模拟数据的基准,这日益耗尽RAM比较reshapedcastdcast.data.table,和spread。我查看了具有三列的简单数据集,具有不同数量的行以获得不同的大小(请参阅最底部的代码)。

> head(df1, 3)
  id                 tms         y
1  1 1970-01-01 01:00:01 0.7463622
2  2 1970-01-01 01:00:01 0.1417795
3  3 1970-01-01 01:00:01 0.6993089
Run Code Online (Sandbox Code Playgroud)

RAM大小仅为8 GB,这是我模拟“非常大”数据集的阈值。为了使计算时间合理,我对每种方法仅进行了3次测量,并专注于从长到宽的重塑。

结果

unit: seconds
       expr       min        lq      mean    median        uq       max neval size.gb …
Run Code Online (Sandbox Code Playgroud)

performance r bigdata reshape

33
推荐指数
1
解决办法
711
查看次数

有没有办法在elasticsearch服务器中导入json文件(包含100个文档).

有没有办法在elasticsearch服务器中导入JSON文件(包含100个文档)?我想将一个大的json文件导入es-server ..

json artificial-intelligence bigdata elasticsearch elasticsearch-plugin

32
推荐指数
6
解决办法
5万
查看次数

sklearn和大型数据集

我有一个22 GB的数据集.我想在我的笔记本电脑上处理它.当然我无法将其加载到内存中.

我使用很多sklearn但是对于更小的数据集.

在这种情况下,经典方法应该是这样的.

只读部分数据 - >部分训练您的估算器 - >删除数据 - >读取其他部分数据 - >继续训练您的估算器.

我已经看到一些sklearn算法具有部分拟合方法,应该允许我们使用数据的各种子样本来训练估计器.

现在我想知道为什么在sklearn中这样做很容易?我正在寻找类似的东西

r = read_part_of_data('data.csv')
m = sk.my_model
`for i in range(n):
     x = r.read_next_chunk(20 lines)
     m.partial_fit(x)

m.predict(new_x)
Run Code Online (Sandbox Code Playgroud)

也许sklearn不是这类东西的正确工具?让我知道.

python bigdata scikit-learn

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

如何在Kafka中使用多个消费者?

我是一名学习卡夫卡的新生,我遇到了一些基本问题,理解了多个消费者,文章,文件等对目前来说并没有太大的帮助.

我试图做的一件事是编写我自己的高级Kafka生产者和消费者并同时运行它们,向主题发布100条简单消息并让我的消费者检索它们.我已经成功地做到了这一点,但是当我尝试引入第二个消费者来消费刚刚发布消息的同一主题时,它不会收到任何消息.

我的理解是,对于每个主题,您可以拥有来自不同消费者群体的消费者,并且每个消费者群体都可以获得针对某个主题生成的消息的完整副本.它是否正确?如果没有,那么建立多个消费者的正确方法是什么?这是我到目前为止写的消费者类:

public class AlternateConsumer extends Thread {
    private final KafkaConsumer<Integer, String> consumer;
    private final String topic;
    private final Boolean isAsync = false;

    public AlternateConsumer(String topic, String consumerGroup) {
        Properties properties = new Properties();
        properties.put("bootstrap.servers", "localhost:9092");
        properties.put("group.id", consumerGroup);
        properties.put("partition.assignment.strategy", "roundrobin");
        properties.put("enable.auto.commit", "true");
        properties.put("auto.commit.interval.ms", "1000");
        properties.put("session.timeout.ms", "30000");
        properties.put("key.deserializer", "org.apache.kafka.common.serialization.IntegerDeserializer");
        properties.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        consumer = new KafkaConsumer<Integer, String>(properties);
        consumer.subscribe(topic);
        this.topic = topic;
    }


    public void run() {
        while (true) {
            ConsumerRecords<Integer, String> records = consumer.poll(0);
            for (ConsumerRecord<Integer, String> record : records) { …
Run Code Online (Sandbox Code Playgroud)

java bigdata apache-kafka

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

当Spark意识到它不再被使用时,Spark会不会自己解决它?

当我们想要多次使用它时,我们可以将RDD持久存储到内存和/或磁盘中.但是,我们以后必须自己解除它们,或者Spark是否会进行某种垃圾收集并在不再需要RDD时解除它的作用?我注意到如果我自己调用unpersist函数,我的性能会变慢.

hadoop distributed-computing bigdata apache-spark rdd

29
推荐指数
1
解决办法
6332
查看次数

Spark镶木地板分区:大量文件

我正在尝试利用spark分区.我试图做类似的事情

data.write.partitionBy("key").parquet("/location")
Run Code Online (Sandbox Code Playgroud)

这里的问题每个分区都会产生大量的镶木地板文件,如果我尝试从根目录中读取,会导致读取速度慢.

为了避免我试过

data.coalese(numPart).write.partitionBy("key").parquet("/location")
Run Code Online (Sandbox Code Playgroud)

但是,这会在每个分区中创建numPart数量的镶木地板文件.现在我的分区大小不同了.所以我理想的是希望每个分区有单独的合并.然而,这看起来并不容易.我需要访问所有分区合并到一定数量并存储在一个单独的位置.

写入后我应该如何使用分区来避免许多文件?

bigdata apache-spark rdd spark-dataframe apache-spark-2.0

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

Cassandra冻结关键词含义

frozenCassandra中关键字的含义是什么?

我正在尝试阅读此文档页面:使用用户定义的类型,但他们对frozen关键字的解释(他们在他们的示例中使用)对我来说不够清楚:

要支持将来的功能,用户定义或元组类型的列定义需要冻结关键字.Cassandra将具有多个组件的冻结值序列化为单个值.有关示例和用法信息,请参阅"使用用户定义的类型","元组类型"和集合类型.

我没有在网上找到任何其他定义或明确的解释.

database keyword bigdata cassandra nosql

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

MapReduce还是Spark?

我用cloudera测试了hadoop和mapreduce,我发现它非常酷,我认为我是最新的相关BigData解决方案.但是几天前,我发现了这个:https: //spark.incubator.apache.org/

一个"闪电快速集群计算系统",能够在Hadoop集群的顶部工作,并且显然能够破坏mapreduce.我看到它在RAM中比mapreduce更有效.我认为当你必须进行集群计算来克服单个机器上的I/O问题时,mapreduce仍然是相关的.但是,由于Spark可以完成mapreduce所做的工作,并且可能在几个操作上更有效率,它不是MapReduce的结束吗?或者MapReduce可以做些什么,或者MapReduce在特定环境中比Spark更有效?

hadoop mapreduce bigdata apache-spark

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