由于数据集较大,无法在此处给出完整的代码。以下是分析步骤和相关代码片段:

  1. 获取数据集

可以从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')
  1. 数据清洗和准备

数据清洗和准备包括去重、缺失值处理、数据类型转换等。可以使用以下代码完成数据清洗和准备:

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')
  1. 数据聚合和分析

根据需求,可以对数据进行聚合和分析。以下是一些示例代码:

# 每个国家/地区的总确诊、死亡和康复病例
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 著作权归作者所有。请勿转载和采集!

免费AI点我,无需注册和登录