2020年美国新冠肺炎疫情数据分析:基于Spark的代码示例
由于数据集较大,无法在此处给出完整的代码。以下是分析步骤和相关代码片段:
- 获取数据集
可以从Johns Hopkins University提供的数据源获取COVID-19数据集。数据集包括每日确诊、死亡和康复病例的国家/地区级别数据。可以使用以下代码从数据源获取数据集:
from pyspark.sql.functions import col, sum
confirmed_cases = spark.read \
.option('header', True) \
.option('inferSchema', True) \
.csv('https://raw.githubusercontent.com/databricks/learning-spark-v2/master/data/Covid/covid_19_data.csv') \
.withColumnRenamed('ObservationDate', 'Date') \
.withColumnRenamed('Country/Region', 'Country') \
.withColumnRenamed('Province/State', 'State') \
.withColumn('Date', col('Date').cast('date'))
deaths = spark.read \
.option('header', True) \
.option('inferSchema', True) \
.csv('https://raw.githubusercontent.com/databricks/learning-spark-v2/master/data/Covid/time_series_covid_19_deaths.csv') \
.withColumnRenamed('Country/Region', 'Country') \
.withColumnRenamed('Province/State', 'State')
recovered = spark.read \
.option('header', True) \
.option('inferSchema', True) \
.csv('https://raw.githubusercontent.com/databricks/learning-spark-v2/master/data/Covid/time_series_covid_19_recovered.csv') \
.withColumnRenamed('Country/Region', 'Country') \
.withColumnRenamed('Province/State', 'State')
- 数据清洗和准备
数据清洗和准备包括去重、缺失值处理、数据类型转换等。可以使用以下代码完成数据清洗和准备:
confirmed_cases = confirmed_cases.dropDuplicates(['Date', 'Country', 'State'])
deaths = deaths.dropDuplicates(['Date', 'Country', 'State'])
recovered = recovered.dropDuplicates(['Date', 'Country', 'State'])
confirmed_cases = confirmed_cases.withColumn('State', col('State').isNull(), 'Unknown')
deaths = deaths.withColumn('State', col('State').isNull(), 'Unknown')
recovered = recovered.withColumn('State', col('State').isNull(), 'Unknown')
confirmed_cases = confirmed_cases.withColumn('Cases', col('Confirmed')) \
.drop('Confirmed')
deaths = deaths.withColumn('Deaths', col('Value')) \
.drop('Value')
recovered = recovered.withColumn('Recovered', col('Value')) \
.drop('Value')
- 数据聚合和分析
根据需求,可以对数据进行聚合和分析。以下是一些示例代码:
# 每个国家/地区的总确诊、死亡和康复病例
total_cases = confirmed_cases.groupBy('Country') \
.agg(sum('Cases').alias('TotalCases'), sum('Deaths').alias('TotalDeaths'), sum('Recovered').alias('TotalRecovered'))
# 每个国家/地区的每日新增确诊、死亡和康复病例
daily_cases = confirmed_cases.groupBy('Country', 'Date') \
.agg(sum('Cases').alias('NewCases'))
daily_deaths = deaths.groupBy('Country', 'Date') \
.agg(sum('Deaths').alias('NewDeaths'))
daily_recovered = recovered.groupBy('Country', 'Date') \
.agg(sum('Recovered').alias('NewRecovered'))
# 以每百万人口为单位计算的总确诊、死亡和康复病例
population = spark.read \
.option('header', True) \
.option('inferSchema', True) \
.csv('https://raw.githubusercontent.com/databricks/learning-spark-v2/master/data/Covid/population.csv')
total_cases_per_million = total_cases.join(population, 'Country') \
.withColumn('CasesPerMillion', col('TotalCases') / col('Population') * 1000000) \
.withColumn('DeathsPerMillion', col('TotalDeaths') / col('Population') * 1000000) \
.withColumn('RecoveredPerMillion', col('TotalRecovered') / col('Population') * 1000000)
原文地址: https://www.cveoy.top/t/topic/oRnt 著作权归作者所有。请勿转载和采集!