【问题标题】:Filtering Spatial Data in Apache Spark在 Apache Spark 中过滤空间数据
【发布时间】:2019-03-05 17:55:05
【问题描述】:

我目前正在解决一个涉及公交车 GPS 数据的问题。我面临的问题是减少我的过程中的计算。

一张表中有大约 20 亿个 GPS 坐标点(经纬度),另一张表中有大约 12,000 个巴士站及其经纬度。预计 20 亿个点中只有 5-10% 位于公交车站。

问题:我只需要标记和提取那些在公共汽车站(12,000 个点)的点(从 20 亿个点中)。由于这是 GPS 数据,因此我无法对坐标进行精确匹配,而是进行基于容差的地理围栏。

问题:在当前的幼稚方法下,标记公交车站的过程需要很长时间。目前,我们正在对 12000 个公交车站点中的每一个点进行挑选,并以 100m 的容差(通过度差转换为距离)查询这 20 亿个点。

问题:是否有算法有效的流程来实现这种点标记?

【问题讨论】:

  • 使用 k-d 树将是开始的地方。
  • 我研究过一个类似的用例。我们使用GeoHashes 的属性来定义单元格并改为定义每个单元格的进程。这仍然是一个广泛的问题。也许您可以展示您当前推动讨论的方法的代码?
  • @LostInOverflow - 当然,通过它。
  • @maasg - 这似乎是个好主意 - 我会试一试!目前,代码只是一组蜂巢查询。但最后的工作需要在 Spark 中完成 - 因此问题就来了。

标签: algorithm apache-spark gps apache-spark-sql geospatial


【解决方案1】:

是的,您可以使用 SpatialSpark 之类的东西。它仅适用于 Spark 1.6.1,但您可以使用 BroadcastSpatialJoin 创建一个非常高效的 RTree

这是我使用 SpatialSpark 和 PySpark 来检查不同多边形是否在彼此内部或相交的示例:

from ast import literal_eval as make_tuple
print "Java Spark context version:", sc._jsc.version()
spatialspark = sc._jvm.spatialspark

rectangleA = Polygon([(0, 0), (0, 10), (10, 10), (10, 0)])
rectangleB = Polygon([(-4, -4), (-4, 4), (4, 4), (4, -4)])
rectangleC = Polygon([(7, 7), (7, 8), (8, 8), (8, 7)])
pointD = Point((-1, -1))

def geomABWithId():
  return sc.parallelize([
    (0L, rectangleA.wkt),
    (1L, rectangleB.wkt)
  ])

def geomCWithId():
  return sc.parallelize([
    (0L, rectangleC.wkt)
  ])

def geomABCWithId():
  return sc.parallelize([
  (0L, rectangleA.wkt),
  (1L, rectangleB.wkt),
  (2L, rectangleC.wkt)])

def geomDWithId():
  return sc.parallelize([
    (0L, pointD.wkt)
  ])

dfAB                 = sqlContext.createDataFrame(geomABWithId(), ['id', 'wkt'])
dfABC                = sqlContext.createDataFrame(geomABCWithId(), ['id', 'wkt'])
dfC                  = sqlContext.createDataFrame(geomCWithId(), ['id', 'wkt'])
dfD                  = sqlContext.createDataFrame(geomDWithId(), ['id', 'wkt'])

# Supported Operators: Within, WithinD, Contains, Intersects, Overlaps, NearestD
SpatialOperator      = spatialspark.operator.SpatialOperator 
BroadcastSpatialJoin = spatialspark.join.BroadcastSpatialJoin

joinRDD = BroadcastSpatialJoin.apply(sc._jsc, dfABC._jdf, dfAB._jdf, SpatialOperator.Within(), 0.0)

joinRDD.count()

results = joinRDD.collect()
map(lambda result: make_tuple(result.toString()), results)

# [(0, 0), (1, 1), (2, 0)] read as:
# ID 0 is within 0
# ID 1 is within 1
# ID 2 is within 0

注意这一行

joinRDD = BroadcastSpatialJoin.apply(sc._jsc, dfABC._jdf, dfAB._jdf, SpatialOperator.Within(), 0.0)

最后一个参数是一个缓冲区值,在您的情况下,它将是您要使用的容差。如果您使用纬度/经度,这可能是一个非常小的数字,因为它是一个径向系统,并且取决于您想要的公差米,您需要calculate based on lat/lon for your area of interest

【讨论】:

    猜你喜欢
    • 2019-04-21
    • 1970-01-01
    • 1970-01-01
    • 2019-03-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-09-05
    • 1970-01-01
    相关资源
    最近更新 更多