### USCRN Data: High-Octane Scraping

In [1]:
import requests
import pandas as pd 
import numpy as np
import yaml 
import re
import itertools
from datetime import datetime
from bs4 import BeautifulSoup

with open ("sources.yaml", "r") as yaml_file:
  sources = yaml.load(yaml_file, Loader=yaml.FullLoader)

##### 1.) Scrape Column Headers and Descriptions 

In [3]:
header_url = sources['USCRN']['headers']
header_response = requests.get(header_url)
header_soup = BeautifulSoup(header_response.content, "html.parser")

columns = str(header_soup).split("\n")[1].strip(" ").split(" ")
columns = list(map(lambda x: str.lower(x), columns)) # columns = [str.lower(c) for c in columns] -- faster?
columns.insert(0,'station_location')

descrip_text = str(header_soup).split("\n")[2] # raw text block containing column descriptions
descrip_text

"The station WBAN number. The UTC date of the observation. The UTC time of the observation. Time is the end of the observed hour, so the 0000 hour is actually the last hour of the previous day's observation (starting just after 11:00 PM through midnight). The Local Standard Time (LST) date of the observation. The Local Standard Time (LST) time of the observation. Time is the end of the observed hour (see UTC_TIME description). The version number of the station datalogger program that was in effect at the time of the observation. Note: This field should be treated as text (i.e. string). Station longitude, using WGS-84. Station latitude, using WGS-84. Average air temperature, in degrees C, during the last 5 minutes of the hour. See Note F. Average air temperature, in degrees C, for the entire hour. See Note F. Maximum air temperature, in degrees C, during the hour. See Note F. Minimum air temperature, in degrees C, during the hour. See Note F. Total amount of precipitation, in mm, record

The descriptions of the columns are quite the mess, as there is no standard separator used. We will have to work our way through it step by step: 

In [4]:
first_split = [re.sub(r'(\([^)]*)$', r"\1)", s) for s in descrip_text.split("). ")] # add ')' back after splitting text on ').' 
no_notes = [re.sub(r' See Note [A-Z]\.',"",s) for s in first_split] # drop any references to notes

The third entry in `no_notes` is ready. The last set of descriptions in `no_notes` can be split on `". "`, but the first two sets need special attention. We will pop the last set out and split it, then pop the third set out, and then address the first two sets. At that point we will recombine everything into one list while preserving the original order. 

In [5]:
last_set = no_notes.pop().strip().split(". ")
third_set = no_notes.pop() # Note: just a string

In [6]:
def flatten(ls:list): 
  return list(itertools.chain.from_iterable(ls)) 

no_notes = [re.sub(". Time is", " at", s) for s in no_notes]
first_second = flatten([s.split(". ") for s in no_notes])

# Finally:
descriptions = flatten([["Location name for USCRN station"], first_second, [third_set], last_set]) # Description added for "station_location" 

In [7]:
header_info = {
  'col_name': columns,
  'description': descriptions, 
  'units': ["X...(Various Lengths)","XXXXX", "YYYYMMDD", "HHmm", "YYYYMMDD", "HHmm", "XXXXXX", "Decimal_degrees", "Decimal_degrees", "Celsius", "Celsius", "Celsius", "Celsius", "mm", "W/m^2", "X", "W/m^2", "X", "W/m^2", "X", "X", "Celsius", "X", "Celsius", "X", "Celsius", "X", "%", "X", "m^3/m^3", "m^3/m^3", "m^3/m^3", "m^3/m^3", "m^3/m^3", "Celsius", "Celsius", "Celsius", "Celsius", "Celsius"]
}

header_df = pd.DataFrame(header_info)
# header_df.to_csv("data/column_descriptions.csv", index=False)

##### 2.) Scrape Core Data Files (>2 million rows)

In [4]:
base_url = sources["USCRN"]["index"]
base_soup = BeautifulSoup(requests.get(base_url).content, "html.parser")

In [5]:
links = base_soup.find_all("a") # 'links' in this notebook will refer to <a> elements, not urls
years = [str(x).zfill(1) for x in range(2000,2024)]
year_links = [link for link in links if link['href'].rstrip('/') in years]

file_urls = []
for year_link in year_links: 
  year_url = base_url + year_link.get("href")
  response = requests.get(year_url) 
  soup = BeautifulSoup(response.content, 'html.parser')
  file_links = soup.find_all('a', href=re.compile(r'AK.*\.txt'))
  if file_links:
    new_file_urls = [year_url + link.getText() for link in file_links]
    file_urls.extend(new_file_urls)

In [6]:
rows = []
regex = r"([St.]*[A-Z][a-z]+_*[A-Za-z]*).*.txt" 
for url in file_urls:
  # Get location from url -- will add to BS results in next step
  file_name = re.search(regex, url).group(0)
  station_location = re.sub("(_formerly_Barrow.*|_[0-9].*)", "", file_name)
  # Get results 
  response = requests.get(url)
  soup = BeautifulSoup(response.content,'html.parser')
  soup_lines = [station_location + " " + line for line in str(soup).strip().split("\n")]
  new_rows = [re.split('\s+', row) for row in soup_lines]
  # Add to list
  rows.extend(new_rows)

In [7]:
df = pd.DataFrame(rows, columns=columns) 

(This dataframe is huge and keeps crashing the kernel when I try to work with it. Save it as a .csv first, restart your kernel, and read it back in before continuing).

In [6]:
# df.to_csv("data/uscrn.csv", index=False)
df = pd.read_csv("data/uscrn.csv")

_From the original data source [README](https://www.ncei.noaa.gov/pub/data/uscrn/products/hourly02/readme.txt):_  

_"Missing data are indicated by the lowest possible integer for a given column format, such as -9999.0 for 7-character fields with one decimal place or -99.000 for 7-character fields with three decimal places."_

We can find these missing value indicators by getting the min of each column.

In [7]:
def minMap(df):
    min_values = {}
    for col in df.columns:
        mv = df[col].min()
        min_values[col] = mv
    return min_values

print(minMap(df))

{'station_location': 'Aleknagik', 'wbanno': 23583, 'utc_date': 20020809, 'utc_time': 0, 'lst_date': 20020808, 'lst_time': 0, 'crx_vn': -9.0, 'longitude': -170.21, 'latitude': 55.05, 't_calc': -9999.0, 't_hr_avg': -9999.0, 't_max': -9999.0, 't_min': -9999.0, 'p_calc': -9999.0, 'solarad': -99999, 'solarad_flag': 0, 'solarad_max': -99999, 'solarad_max_flag': 0, 'solarad_min': -99999, 'solarad_min_flag': 0, 'sur_temp_type': 'C', 'sur_temp': -9999.0, 'sur_temp_flag': 0, 'sur_temp_max': -9999.0, 'sur_temp_max_flag': 0, 'sur_temp_min': -9999.0, 'sur_temp_min_flag': 0, 'rh_hr_avg': -9999, 'rh_hr_avg_flag': 0, 'soil_moisture_5': -99.0, 'soil_moisture_10': -99.0, 'soil_moisture_20': -99.0, 'soil_moisture_50': -99.0, 'soil_moisture_100': -99.0, 'soil_temp_5': -9999.0, 'soil_temp_10': -9999.0, 'soil_temp_20': -9999.0, 'soil_temp_50': -9999.0, 'soil_temp_100': -9999.0}


We will replace these values with `NaNs`, but we need to be careful: since the source does not normally have empty records, any `NaNs` entering our pipeline on read will likely come either from errors in the data source or errors in our attempts to read from it. When writing our update DAG, before we replace any values with `NaNs` we'll need to check for `NaNs` and log an alert if any are found. 

In [8]:
df.replace([-99999,-9999], np.nan, inplace=True) # Can safely assume these are always missing values in every column they appear in
df = df.filter(regex="^((?!soil).)*$") # vast majority of soil columns have missing data
df.replace({'crx_vn':{-9:np.nan}}, inplace=True)

Next, let's convert the date and time columns to `datetime` objects and reorder our columns

In [9]:
df['utc_datetime'] = pd.to_datetime(df['utc_date'].astype(int).astype(str) + df['utc_time'].astype(int).astype(str).str.zfill(4), format='%Y%m%d%H%M')
df['lst_datetime'] = pd.to_datetime(df['lst_date'].astype(int).astype(str) + df['lst_time'].astype(int).astype(str).str.zfill(4), format='%Y%m%d%H%M')

In [10]:
# drop old date and time columns
df.drop(['utc_date', 'utc_time', 'lst_date', 'lst_time'], axis=1, inplace=True)
df.columns

Index(['station_location', 'wbanno', 'crx_vn', 'longitude', 'latitude',
       't_calc', 't_hr_avg', 't_max', 't_min', 'p_calc', 'solarad',
       'solarad_flag', 'solarad_max', 'solarad_max_flag', 'solarad_min',
       'solarad_min_flag', 'sur_temp_type', 'sur_temp', 'sur_temp_flag',
       'sur_temp_max', 'sur_temp_max_flag', 'sur_temp_min',
       'sur_temp_min_flag', 'rh_hr_avg', 'rh_hr_avg_flag', 'utc_datetime',
       'lst_datetime'],
      dtype='object')

In [13]:
# reorder columns 
cols = ['station_location','wbanno','crx_vn','utc_datetime','lst_datetime'] + list(df.columns)[3:-2]
df = df[cols]

Lastly, let's add a `date_added` column: 

In [15]:
df['date_added_utc'] = datetime.now()

In [16]:
df.sample(5)

Unnamed: 0,station_location,wbanno,crx_vn,utc_datetime,lst_datetime,longitude,latitude,t_calc,t_hr_avg,t_max,...,sur_temp_type,sur_temp,sur_temp_flag,sur_temp_max,sur_temp_max_flag,sur_temp_min,sur_temp_min_flag,rh_hr_avg,rh_hr_avg_flag,date_added_utc
1638727,Port_Alsworth,26562.0,2.424,2020-02-15 12:00:00,2020-02-15 03:00:00,-154.32,60.2,-11.3,-11.0,-10.8,...,C,-11.9,0.0,-11.0,0.0,-12.8,0.0,84.0,0.0,2023-02-17 16:10:29.402040
1988997,Glennallen,56401.0,2.515,2022-01-30 02:00:00,2022-01-29 17:00:00,-145.5,63.03,-12.8,-12.4,-12.1,...,C,-16.5,0.0,-15.8,0.0,-17.2,0.0,72.0,0.0,2023-02-17 16:10:29.402040
1067823,Port_Alsworth,26562.0,2.424,2017-12-13 22:00:00,2017-12-13 13:00:00,-154.32,60.2,2.7,3.0,4.1,...,C,0.1,0.0,0.1,0.0,0.0,0.0,91.0,0.0,2023-02-17 16:10:29.402040
1385979,Glennallen,56401.0,2.515,2019-04-22 08:00:00,2019-04-21 23:00:00,-145.5,63.03,-0.5,-0.8,-0.1,...,C,-4.7,0.0,-4.5,0.0,-4.9,0.0,71.0,0.0,2023-02-17 16:10:29.402040
1166338,Deadhorse,26565.0,2.514,2018-06-17 18:00:00,2018-06-17 09:00:00,-148.46,70.16,-0.3,-0.5,-0.1,...,C,3.9,0.0,4.8,0.0,3.2,0.0,88.0,0.0,2023-02-17 16:10:29.402040


In [17]:
df.to_csv("data/uscrn.csv", index=False)

Let's also make a table for our various station locations. This will be useful when searching for the four-day forecasts in the NWS notebook. 

In [1]:
locations = df[['station_location', 'wbanno', 'longitude', 'latitude']].drop_duplicates()
# locations.to_csv("data/locations.csv", index=False)

##### 3.) Upload Core Data to BigQuery 

In [3]:
# df=pd.read_csv("data/uscrn.csv") <-- restarted kernel again before here

In [8]:
df.head()

Unnamed: 0,station_location,wbanno,crx_vn,utc_datetime,lst_datetime,longitude,latitude,t_calc,t_hr_avg,t_max,...,sur_temp_type,sur_temp,sur_temp_flag,sur_temp_max,sur_temp_max_flag,sur_temp_min,sur_temp_min_flag,rh_hr_avg,rh_hr_avg_flag,date_added_utc
0,Fairbanks,26494.0,1.001,2002-08-09 22:00:00.000000,2002-08-09 13:00:00.000000,-147.51,64.97,,,,...,R,19.4,3.0,,0.0,,0.0,0.0,3.0,2023-02-17 16:10:29.402040
1,Fairbanks,26494.0,1.001,2002-08-09 23:00:00.000000,2002-08-09 14:00:00.000000,-147.51,64.97,10.0,11.1,12.0,...,R,14.6,0.0,,0.0,,0.0,0.0,3.0,2023-02-17 16:10:29.402040
2,Fairbanks,26494.0,1.001,2002-08-10 00:00:00.000000,2002-08-09 15:00:00.000000,-147.51,64.97,12.0,10.5,12.0,...,R,15.0,0.0,,0.0,,0.0,0.0,3.0,2023-02-17 16:10:29.402040
3,Fairbanks,26494.0,1.001,2002-08-10 01:00:00.000000,2002-08-09 16:00:00.000000,-147.51,64.97,12.1,11.9,12.2,...,R,15.4,0.0,,0.0,,0.0,0.0,3.0,2023-02-17 16:10:29.402040
4,Fairbanks,26494.0,1.001,2002-08-10 02:00:00.000000,2002-08-09 17:00:00.000000,-147.51,64.97,11.9,12.0,12.1,...,R,14.2,0.0,,0.0,,0.0,0.0,3.0,2023-02-17 16:10:29.402040


In [79]:
%%bash
bq mk -d --location=us-east4 team-week3:alaska

Dataset 'team-week3:alaska' successfully created.


Core Data: 

In [12]:
from google.cloud import bigquery
from google.oauth2 import service_account

# Setting certain numeric columns (e.g. crx_vn, the flag columns) as strings will indicate that they are not meant to have arithmetic calculations done on them
schema = [
  bigquery.SchemaField("station_location", "STRING", mode="REQUIRED"), 
  bigquery.SchemaField("wbanno", "STRING", mode="REQUIRED"), 
  bigquery.SchemaField("crx_vn", "STRING", mode="NULLABLE"), 
  bigquery.SchemaField("utc_datetime", "DATETIME", mode="REQUIRED"), 
  bigquery.SchemaField("lst_datetime", "DATETIME", mode="REQUIRED"), 
  bigquery.SchemaField("longitude", "FLOAT", mode="REQUIRED"), 
  bigquery.SchemaField("latitude", "FLOAT", mode="REQUIRED"), 
  bigquery.SchemaField("t_calc", "FLOAT", mode="NULLABLE"), 
  bigquery.SchemaField("t_hr_avg", "FLOAT", mode="NULLABLE"), 
  bigquery.SchemaField("t_max", "FLOAT", mode="NULLABLE"), 
  bigquery.SchemaField("t_min", "FLOAT", mode="NULLABLE"), 
  bigquery.SchemaField("p_calc", "FLOAT", mode="NULLABLE"), 
  bigquery.SchemaField("solarad", "FLOAT", mode="NULLABLE"), 
  bigquery.SchemaField("solarad_flag", "STRING", mode="NULLABLE"), 
  bigquery.SchemaField("solarad_max", "FLOAT", mode="NULLABLE"), 
  bigquery.SchemaField("solarad_max_flag", "STRING", mode="NULLABLE"), 
  bigquery.SchemaField("solarad_min", "FLOAT", mode="NULLABLE"), 
  bigquery.SchemaField("solarad_min_flag", "STRING", mode="NULLABLE"), 
  bigquery.SchemaField("sur_temp_type", "STRING", mode="NULLABLE"), 
  bigquery.SchemaField("sur_temp", "FLOAT", mode="NULLABLE"), 
  bigquery.SchemaField("sur_temp_flag", "STRING", mode="NULLABLE"), 
  bigquery.SchemaField("sur_temp_max", "FLOAT", mode="NULLABLE"), 
  bigquery.SchemaField("sur_temp_max_flag", "STRING", mode="NULLABLE"), 
  bigquery.SchemaField("sur_temp_min", "FLOAT", mode="NULLABLE"), 
  bigquery.SchemaField("sur_temp_min_flag", "STRING", mode="NULLABLE"), 
  bigquery.SchemaField("rh_hr_avg", "FLOAT", mode="NULLABLE"), 
  bigquery.SchemaField("rh_hr_avg_flag", "STRING", mode="NULLABLE"), 
  bigquery.SchemaField("date_added_utc", "DATETIME", mode="REQUIRED")
]

In [13]:
key_path = "/home/alex/.creds/alex-sa-tw3.json"
credentials = service_account.Credentials.from_service_account_file(
   key_path, scopes=["https://www.googleapis.com/auth/cloud-platform"],
)

client = bigquery.Client(credentials=credentials, project=credentials.project_id)

table_id = f"{credentials.project_id}.alaska.uscrn"

jc = bigquery.LoadJobConfig(
   source_format = bigquery.SourceFormat.CSV,
   autodetect=False,
   schema=schema,
   create_disposition="CREATE_IF_NEEDED",
   write_disposition="WRITE_TRUNCATE", 
   destination_table_description="Historical weather data from USCRN stations in Alaska"
)

job = client.load_table_from_dataframe(df, table_id, job_config=jc)

job.result()

LoadJob<project=team-week3, location=us-east4, id=8dbb2c55-bcc8-491d-8c23-e24c563dc7a6>

##### 4.) DAG Task for Updating Dataset  

Pandas has a neat function for reading HTML tables to dataframes (`pd.read_html`). It's not ideal for "messier" tabular data or for iterating through lots of HTML pages like we did earlier (iteratively creating and appending dataframes is very slow given the size of dataframe objects). But it's definitely useful for reading a single table object:   

In [10]:
now = datetime.now()
updates_url = sources['USCRN']['updates'] + f"{now.year}"

df = pd.read_html(updates_url, skiprows=[1,2])[0]
df.drop(["Size", "Description"], axis=1, inplace=True)
df.dropna(inplace=True)
cols = [re.sub(" ","_",str.lower(c)) for c in df.columns]
df.columns = cols
df['last_modified'] = pd.to_datetime(df['last_modified'])
df

Unnamed: 0,name,last_modified
0,CRN60H0203-202301010100.txt,2022-12-31 20:47:00
1,CRN60H0203-202301010200.txt,2022-12-31 21:47:00
2,CRN60H0203-202301010300.txt,2022-12-31 22:53:00
3,CRN60H0203-202301010400.txt,2022-12-31 23:54:00
4,CRN60H0203-202301010500.txt,2023-01-01 00:49:00
...,...,...
1138,CRN60H0203-202302171100.txt,2023-02-17 06:47:00
1139,CRN60H0203-202302171200.txt,2023-02-17 07:47:00
1140,CRN60H0203-202302171300.txt,2023-02-17 08:47:00
1141,CRN60H0203-202302171400.txt,2023-02-17 09:47:00


In [21]:
from collections import deque
from io import StringIO

with open("data/uscrn.csv", 'r') as fp:
    q = deque(fp, 1)  
last_added = pd.read_csv(StringIO(''.join(q)), header=None).iloc[0,-1]
last_added = datetime.strptime(last_added, "%Y-%m-%d %H:%M:%S.%f")

In [22]:
df[df['last_modified'] > last_added] # In actual use this won't be empty 

Unnamed: 0,name,last_modified


In [None]:
new_file_urls = updates_url + df['name']

rows = []
regex = r"([St.]*[A-Z][a-z]+_*[A-Za-z]*).*.txt" 
for url in new_file_urls:
  # Get location from url -- will add to BS results in next step
  file_name = re.search(regex, url).group(0)
  station_location = re.sub("(_formerly_Barrow.*|_[0-9].*)", "", file_name)
  # Get results 
  response = requests.get(url)
  soup = BeautifulSoup(response.content,'html.parser')
  soup_lines = [station_location + " " + line for line in str(soup).strip().split("\n")]
  new_rows = [re.split('\s+', row) for row in soup_lines]
  # Add to list
  rows.extend(new_rows)