【问题标题】:Large web data sets in Python - how to deal with very big arrays?Python 中的大型 Web 数据集 - 如何处理非常大的数组?
【发布时间】:2016-05-14 08:28:43
【问题描述】:

我们目前正在使用 Cassandra (http://cassandra.apache.org/) 处理时间序列数据。 Cassandra 的读取速度非常快,但我们必须在呈现数据之前对数据执行一系列计算(实际上我们是在模仿 SQL 的 SUM 和 GROUP BY 功能 - Cassandra 不支持开箱即用的功能)

我们(在一定程度上)熟悉 Python,并决定构建一个脚本来查询我们的 Cassandra 集群以及执行数学运算并以 JSON 格式呈现结果:

query = (
    "SELECT query here...")

startTimeQuery = time.time()

# Executes cassandra query
rslt = cassession.execute(query)

print("--- %s seconds to query ---" % (time.time() - startTimeQuery))

tally = {}

startTimeCalcs = time.time()
for row in rslt:
    userid = row.site_user_id

    revenue = (int(row.revenue) - int(row.reversals_revenue or 0))
    accepted = int(row.accepted or 0)
    reversals_revenue = int(row.reversals_revenue or 0)
    error = int(row.error or 0)
    impressions_negative = int(row.impressions_negative or 0)
    impressions_positive = int(row.impressions_positive or 0)
    rejected = int(row.rejected or 0)
    reversals_rejected = int(row.reversals_rejected or 0)

    if tally.has_key(userid):
        tally[userid]["revenue"] += revenue
        tally[userid]["accepted"] += accepted
        tally[userid]["reversals_revenue"] += reversals_revenue
        tally[userid]["error"] += error
        tally[userid]["impressions_negative"] += impressions_negative
        tally[userid]["impressions_positive"] += impressions_positive
        tally[userid]["rejected"] += rejected
        tally[userid]["reversals_rejected"] += reversals_rejected
    else:
        tally[userid] = {
            "accepted": accepted,
            "error": error,
            "impressions_negative": impressions_negative,
            "impressions_positive": impressions_positive,
            "rejected": rejected,
            "revenue": revenue,
            "reversals_rejected": reversals_rejected,
            "reversals_revenue": reversals_revenue
        }


print("--- %s seconds to calculate results ---" % (time.time() - startTimeCalcs))

startTimeJson = time.time()
jsonOutput =json.dumps(tally)
print("--- %s seconds for json dump ---" % (time.time() - startTimeJson))

print("--- %s seconds total ---" % (time.time() - startTimeQuery))

print "Array Size: " + str(len(tally)) 

这是我们得到的输出:

--- 0.493520975113 seconds to query ---
--- 23.1472680569 seconds to calculate results ---
--- 0.546246051788 seconds for json dump ---
--- 24.1871240139 seconds total ---
Array Size: 198124

我们在计算上花费了大量时间,我们知道问题不在于求和和分组本身:问题在于数组的绝对大小。

我们听说过一些关于 numpy 的好消息,但我们数据的性质使得矩阵大小成为未知数。

我们正在寻找有关如何解决此问题的任何提示。包括完全不同的编程方法。

【问题讨论】:

  • 时间序列数据的 goto python 包是pandas,它在后台使用numpy。你调查过吗?
  • 还有,“大”有多大?
  • 我们在 Cassandra 的 Python 驱动程序文档中找到以下内容: NumpyProtocolHandler :可用于将结果直接反序列化为 NumPy 数组。这有助于与分析工具包(如 Pandas)的有效集成。不过,我们在查找任何类型的更广泛的文档或示例时遇到了问题。

标签: python numpy cassandra bigdata


【解决方案1】:

我做了一个非常相似的处理,我也担心处理时间。我认为您没有考虑重要的事情:您从 cassandra 收到的结果对象作为函数 execute() 的返回不包含您想要的所有行。相反,它包含一个分页结果,并在您扫过for 列表中的对象时获得行。虽然这是基于个人观察,但我不知道要提供更多技术细节。

我建议您通过在execute 命令之后添加一个简单的rslt = list(rslt) 来隔离结果的查询和处理,这将强制python 在执行处理之前遍历结果中的所有行,同时强制cassandra 驱动程序在进行处理之前获取您想要的所有行。

我想你会发现你的很多处理时间实际上是在查询,但是驱动程序通过分页结果掩盖了它。

【讨论】:

    【解决方案2】:

    Cassandra 2.2 及更高版本允许用户定义聚合函数。您可以使用它在 cassandra 端执行列排序。 请参阅DataStax article 了解有关用户定义聚合的数据

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2013-01-11
      • 2021-10-19
      • 1970-01-01
      • 1970-01-01
      • 2010-10-07
      相关资源
      最近更新 更多