【问题标题】:PySpark equivalent of Pandas statementsPySpark 相当于 Pandas 语句
【发布时间】:2020-10-22 05:56:22
【问题描述】:

我有一个问题。我有一个正在清理的 Spark 数据框,在 pandas 中,我通常会使用:

df['cases_since_prev_day'] = df.groupby(['county','state'])['cases'].diff().fillna(0).astype(int)

我可以在 PySpark 中使用聚合的等效语句吗?

编辑:所以,原始数据集来自 Github 上的这张表 - https://raw.githubusercontent.com/nytimes/covid-19-data/master/us-counties.csv

快速浏览如下:

date,county,state,fips,cases,deaths
2020-01-21,Snohomish,Washington,53061,1,0
2020-01-22,Snohomish,Washington,53061,1,0
2020-01-23,Snohomish,Washington,53061,1,0
2020-01-24,Cook,Illinois,17031,1,0
2020-01-24,Snohomish,Washington,53061,1,0
2020-01-25,Orange,California,06059,1,0
2020-01-25,Cook,Illinois,17031,1,0
2020-01-25,Snohomish,Washington,53061,1,0

我之前用pandas清理过数据集,生成的形式如下:

date,county,state,fips,cases,deaths,ISO3166_1,ISO3166_2,cases_since_prev_day,deaths_since_prev_day,Last_Update_Date,Last_Reported_Flag
2020-03-19,Abbeville,South Carolina,45001,1,0,US,SC,0,0,2020-10-21 22:34:14.644190,False

我想使用 PySpark 完成同样的任务。到目前为止,我有代码和输出:

#Data from The New York Times, based on reports from state and local health agencies
import pandas as pd
import datetime
import pycountry
import numpy as np
import sys
sys.path.append('../utilities')

from utility_setup import create_spark_session, read_s3_to_dataframes
from pyspark.sql.functions import col, create_map, lit
from pyspark.sql.functions import udf
from write_to_s3 import _write_dataframe_to_csv

from pyspark import SparkContext
from pyspark.sql import SQLContext
from pyspark.sql import functions as F
from pyspark.sql.window import Window

#Dataframe - read CSV
bucket = 'covid19datalakesafford'
key = 'us-counties.csv'
spark = create_spark_session('COVID-19 NYT - county cases')
df = read_s3_to_dataframes(spark, bucket, key)
df = df.withColumn("fips",df["fips"].cast('string')).withColumn("ISO3166_1",lit("US"))
#df.state.show(2)

subdivisions = {k.name: k.code.replace("US-", "") for k in pycountry.subdivisions.get(country_code="US")}
mapping_func = lambda x: subdivisions.get(x) 
df = df.withColumn('ISO3166_2',udf(mapping_func)("state"))

#df_2["ISO3166_2"] = df_2["state"].apply(lambda x: subdivisions.get(x))



#Old way using Python alone
#df = pd.read_csv("https://raw.githubusercontent.com/nytimes/covid-19-data/master/us-counties.csv", \
#     dtype={'fips': str})

#df = df.withColumnRenamed("ISO3166_1","US") \
#    .withColumnRenamed("ISO3166_2","state")
#df.show(2)

#subdivision_udf = udf(lambda x: subdivisions.get(x))

#function for applying dictionary terms for subdivisions to column

df = df.sort('county', 'date', 'ISO3166_1', 'ISO3166_2')
df.show(2)

#Equivalent in PySpark for lines below?

#df['cases_since_prev_day'] = df.groupby(['county','state'])['cases'].diff().fillna(0).astype(int)
#df["Last_Update_Date"] = datetime.datetime.utcnow()
#df['Last_Reported_Flag'] = df['date'] == df['date'].max()

当前表(前两行):

+----------+---------+--------------+-----+-----+------+---------+---------+
|      date|   county|         state| fips|cases|deaths|ISO3166_1|ISO3166_2|
+----------+---------+--------------+-----+-----+------+---------+---------+
|2020-03-19|Abbeville|South Carolina|45001|    1|     0|       US|       SC|
|2020-03-20|Abbeville|South Carolina|45001|    1|     0|       US|       SC|
+----------+---------+--------------+-----+-----+------+---------+---------+

编辑 2: 请注意,这是一个时间序列,我预计 COVID-19 病例的列表每天都会增加,因为它们每天都会由每个县/地区报告每个州。目前,我的表中的行数接近 950,000 行,而使用 Pandas 的行数很慢(需要 9 分钟才能完成)。

【问题讨论】:

  • 我们需要更多信息 - 请以表格格式分享输入和输出数据是什么样的
  • @dsk 我已经编辑了帖子。

标签: python apache-spark pyspark


【解决方案1】:

这应该让你(几乎)一直到那里(我没有你的映射函数,所以不能得到两个字母的状态缩写):

from pyspark.sql.window import Window
from pyspark.sql import functions as F

# county-level windowing - like in SQL
win = Window().partitionBy('county', 'state').orderBy('date')

# lags and final date
df = df.withColumn('cases_since_prev_day', F.lag('cases').over(win))
df = df.withColumn('deaths_since_prev_day', F.lag('deaths').over(win))
df = df.withColumn('Last_Update_Date', F.last('date').over(win))

【讨论】:

  • 谢谢!至少对我来说这是一个起点!
  • 如果有帮助,您介意点赞还是接受?
  • 嗨@ags29,你会碰巧知道如何在 PySpark 中这样做: data['pendingIncrease'] = data['pending'] - data.groupby(['Province_State'])[ "待定"].shift(1)
  • 是的,我可以在 PySpark 中编写等价的代码,稍后再做
  • win = Window().partitionBy('Province_State').orderBy('date') data = data.withColumn('pendingIncrease', F.col('pending') - F.lag('pending').over(win)
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-02-04
  • 1970-01-01
  • 2022-12-03
  • 2021-12-23
  • 2018-04-21
  • 2017-05-12
  • 2013-01-29
相关资源
最近更新 更多