## Carbon Monitoring Project

In [1]:
import holoviews as hv
import pandas as pd
hv.extension('bokeh')

This notebook aims to visualize the data used in the carbon monitoring project [nee_data_fusion](https://github.com/greyNearing/nee_data_fusion/) using Python tools.

The goals of this notebook:

* examine the measurements from each site
* generate some visualization or global model to predict one site from every other site.
* generate and explain model idea

To run this notebook, you will need to symlink the `data` directory of the [nee_data_fusion](https://github.com/greyNearing/nee_data_fusion/) to `flux_data` in the `examples` directory of `EarthML`. In addition, you will need `RSIF_2007_2016_05N_01L.mat` in the `examples` directory which you can download from https://gentinelab.eee.columbia.edu/content/datasets

### Loading FluxNet data ``extract_fluxnet.m``

[FluxNet](http://fluxnet.fluxdata.org/) is a worldwide collection of sensor stations that record a number of local variables relating to atmospheric conditions, solar flux and soil moisture. The data is the [nee_data_fusion](https://github.com/greyNearing/nee_data_fusion/) repository is expressed as a collection of CSV files where the site names are expressed in the filenames.

This cell defines functions to

* read in the data from all sites
* do some data munging (i.e., date parsing, `NaN` replacement)

In [2]:
import numpy as np
import datetime
import os

In [3]:
# DAILIES_DIR = '../../nee_data_fusion/data/in_situ/fluxnet_daily/'
# METADATA_CSV = '../../nee_data_fusion/data/in_situ/extracted/allflux_metadata.txt'
DAILIES_DIR = 'flux_data/dailies/'
METADATA_CSV = 'allflux_metadata.txt'

sites = [fname.split('_')[1] for fname in os.listdir(DAILIES_DIR)]
metadata = pd.read_csv(METADATA_CSV, header=None, names=['site', 'lat', 'lon', 'igbp', 'network'], 
                       usecols=['site', 'lat', 'lon', 'igbp'], index_col='site')

# get all the igbp codes for these sites
igbp_codes = metadata.loc[sites].igbp.unique()

Passing list-likes to .loc or [] with any missing label will raise
KeyError in the future, you can use .reindex() as an alternative.

See the documentation here:
http://pandas.pydata.org/pandas-docs/stable/indexing.html#deprecate-loc-reindex-listlike
  # This is added back by InteractiveShellApp.init_path()


In [4]:
# Any missing metadata?
metadata.loc[sites].isnull().values.any()

Passing list-likes to .loc or [] with any missing label will raise
KeyError in the future, you can use .reindex() as an alternative.

See the documentation here:
http://pandas.pydata.org/pandas-docs/stable/indexing.html#deprecate-loc-reindex-listlike
  


True

In [5]:
def _parse_days(integer):
    """ `integer` date as `20180704` to represent July 4th, 2018"""
    x = str(integer)
    d = {'year': int(x[:4]), 'month': int(x[4:6]), 'day': int(x[6:])}
    day_of_year = datetime.datetime(d['year'], d['month'], d['day']).timetuple().tm_yday
    return day_of_year

def clean(df, timestamp_col="TIMESTAMP", site='', keep=[], drop=[], predict=''):
    """
    Clean the dataset
    
    * Replace NaN's and any number less than -9990 with 0s
    * drop columns specified in `drop`
    * pull out prediction and feature matrices
    * Parse timestamp and pull out day of year ("DOY")
    """
    limit = -9990
    for i in range(50):
        df = df.replace(limit - i, np.nan)
    
    to_drop = [col for col in drop if col in df.columns]
    df.drop(columns=to_drop, inplace=True)
    for col in keep:
        if col not in df.columns:
            if 'SWC_F' in col or 'TS_F' in col:
                df[col] = 0
    
    df = df.fillna(0)
    df['DOY'] = df['TIMESTAMP'].apply(_parse_days)  
    df.pop('TIMESTAMP')
    X = df[keep]
    y = df[predict]
    return X, y

def load_fluxnet_site(site, one=False):
    """
    The main function to load data
    
    Parameters
    ----------
    site : str
        e.g., "US-CA1"
    one : bool, optional
        Whether to preform a dirty hack and create "one" dataframe that
        includes the prediction variable in the feature matrix.
        
    Returns
    -------
    X : pd.DataFrame
        Feature matrix. If ``one``, this will include the prediction variable
        and be the only thing returned.
    y : pd.DataFrame
        The prediction variable. Not returned if ``one``.
    """
    #dataRaw(dataRaw <= -9990) = 0/0 (is NaN?)
    #NaN -> zero
    prefix = 'FLX_{site}_FLUXNET'.format(site=site)
    filenames = [fname for fname in os.listdir(DAILIES_DIR)
                if fname.startswith(prefix)]
    if len(filenames) != 1:
        raise FileNotFoundError
    filename = filenames[0]
    path = '{directory}{filename}'.format(directory=DAILIES_DIR, filename=filename)
    
    raw_daily = pd.read_csv(path)    
    
    keep =  ['P_ERA',
             'TA_ERA',
             'PA_ERA',
             'SW_IN_ERA',
             'LW_IN_ERA',
             'WS_ERA',
             'SWC_F_MDS_1', 'SWC_F_MDS_2', 'SWC_F_MDS_3',
             'TS_F_MDS_1', 'TS_F_MDS_2', 'TS_F_MDS_3',
             'VPD_ERA',
             'DOY']
    drop = ["GPP_DT_VUT_USTAR50",
            "GPP_DT_CUT_USTAR50",
            "LE_F_MDS",
            "H_F_MDS"]
    predict = ["NEE_CUT_USTAR50",
               "NEE_VUT_USTAR50"]
    
    X, y = clean(raw_daily, keep=keep, drop=drop, predict=predict[0])
    X['site'] = site  # some metadata
    X['y'] = y
    return X

## Setup Dask
Dask is required to read in the CSVs and do preprocessing *quickly*.

In [6]:
import dask
import dask.array as da
import dask.dataframe as dd
from distributed import Client

client = Client()
client

0,1
Client  Scheduler: tcp://127.0.0.1:51563  Dashboard: http://127.0.0.1:8787/status,Cluster  Workers: 8  Cores: 8  Memory: 17.18 GB


## Read in data

In [7]:
futures = client.map(load_fluxnet_site, sites, one=True)

In [8]:
succeeded = [f for f in futures if not f.exception()]
failed = [f for f in futures if f.exception()]

In [9]:
dfs = client.gather(succeeded)

## Merge data

Once the data are loaded in, they need to be joined with the metadata relating to each site.

In [10]:
df = pd.concat(dfs)
df.columns

Index(['P_ERA', 'TA_ERA', 'PA_ERA', 'SW_IN_ERA', 'LW_IN_ERA', 'WS_ERA',
       'SWC_F_MDS_1', 'SWC_F_MDS_2', 'SWC_F_MDS_3', 'TS_F_MDS_1', 'TS_F_MDS_2',
       'TS_F_MDS_3', 'VPD_ERA', 'DOY', 'site', 'y'],
      dtype='object')

In [11]:
# create a little dataframe for mapping categorical variables

# TODO: it looks like this is one hot encoding. Use pd.get_dummies instead?
# if not, there's another TODO:
# TODO: it looks like this is creating an identity matrix. Use np.eye?
a = np.zeros((len(igbp_codes), len(igbp_codes)), int)
np.fill_diagonal(a, 1)
categorical_igbp_mapper = pd.DataFrame(index=igbp_codes, columns=igbp_codes, data=a)
categorical_igbp_mapper.rename_axis('igbp', inplace=True)

# add metadata to the big dataframe
onehot_metadata = pd.get_dummies(metadata, columns=['igbp'])
onehot_metadata['igbp'] = metadata['igbp']
assert onehot_metadata.index.name == 'site'
onehot_metadata['site'] = onehot_metadata.index

In [12]:
df = pd.merge(df, onehot_metadata, on='site')

In [13]:
# show = df.sample(frac=0.10)
show = df.copy()

In [14]:
# set this to False to not include vegetation type in the calculation
include_veg = True

sites = pd.Categorical(show['site']).codes
dropped = {}
for col in ['DOY', 'site', 'igbp', 'lat', 'lon']:
    dropped[col] = show[col].copy()
    show.pop(col)
    
if not include_veg:
    for col in igbp_codes:
        dropped[col] = show[col].copy()
        show.pop(col)
        
print("{} observations and {} variables".format(*show.shape))
print("Generating a prediction with these variables: \n  {}".format(
    "\n  ".join(list(
        show.columns
    ))
))

532544 observations and 29 variables
Generating a prediction with these variables: 
  P_ERA
  TA_ERA
  PA_ERA
  SW_IN_ERA
  LW_IN_ERA
  WS_ERA
  SWC_F_MDS_1
  SWC_F_MDS_2
  SWC_F_MDS_3
  TS_F_MDS_1
  TS_F_MDS_2
  TS_F_MDS_3
  VPD_ERA
  y
  igbp_BSV
  igbp_CRO
  igbp_CSH
  igbp_DBF
  igbp_DNF
  igbp_EBF
  igbp_ENF
  igbp_GRA
  igbp_MF
  igbp_OSH
  igbp_SAV
  igbp_SNO
  igbp_WAT
  igbp_WET
  igbp_WSA


These variables are sufficient to create the linear models at every site. However, the site information is hidden from the visualization algorithm.

* Good sanity checks:
    - latitude encoded some structure, longitude does not

## Visualization

Linear models work well *at one site* but this is confounded by

* lat/lon
* day of year
* environment type

We want to generate some visualization that accounts for these 4 variables and helps generate some understanding.

That is, these observations lie on some manifold. We want to learn the structure of that manifold, and visualize each observation on that manifold.

This work attempts to find similar observations - observations that have a similar structure between the independent variables (e.g., `P_ERA`) and dependent variables (the carbon flux measurement `y`).

UMAP is a tool for this, and has firm mathematical grounding (plus, it's nice to use).

In [15]:
import umap
reduct = umap.UMAP(verbose=True, n_epochs=None)#, n_neighbors=30)

UMAP(a=None, angular_rp_forest=False, b=None, init='spectral',
   learning_rate=1.0, local_connectivity=1.0, metric='euclidean',
   metric_kwds=None, min_dist=0.1, n_components=2, n_epochs=None,
   n_neighbors=15, negative_sample_rate=5, random_state=None,
   repulsion_strength=1.0, set_op_mix_ratio=1.0, spread=1.0,
   target_metric='categorical', target_metric_kwds=None,
   target_n_neighbors=-1, target_weight=0.5, transform_queue_size=4.0,
   transform_seed=42, verbose=True)


In [16]:
reduct.fit(show.values)

Construct fuzzy simplicial set
	 0  /  16
	 1  /  16
	 2  /  16
Construct embedding
	completed  0  /  200 epochs
	completed  20  /  200 epochs
	completed  40  /  200 epochs
	completed  60  /  200 epochs
	completed  80  /  200 epochs
	completed  100  /  200 epochs
	completed  120  /  200 epochs
	completed  140  /  200 epochs
	completed  160  /  200 epochs
	completed  180  /  200 epochs


UMAP(a=None, angular_rp_forest=False, b=None, init='spectral',
   learning_rate=1.0, local_connectivity=1.0, metric='euclidean',
   metric_kwds=None, min_dist=0.1, n_components=2, n_epochs=None,
   n_neighbors=15, negative_sample_rate=5, random_state=None,
   repulsion_strength=1.0, set_op_mix_ratio=1.0, spread=1.0,
   target_metric='categorical', target_metric_kwds=None,
   target_n_neighbors=-1, target_weight=0.5, transform_queue_size=4.0,
   transform_seed=42, verbose=True)

In [17]:
embedding = reduct.embedding_
embedding

array([[ 2.176026 , -0.7342752],
       [10.628123 ,  1.2661803],
       [ 3.389549 , -0.8920688],
       ...,
       [ 5.6134677, -0.9864353],
       [-5.568358 , -4.9773297],
       [ 6.50224  , -2.4452446]], dtype=float32)

In [18]:
cols = ['lat', 'lon', 'igbp']
s = pd.DataFrame(dropped)
s['x0'] = embedding[:, 0]
s['x1'] = embedding[:, 1]
for col in cols:
    if col in show:
        s[col] = show[col]
    else:
        if not col in s:
            print(col)

In [19]:
from bokeh.models import Select
from bokeh.layouts import row, widgetbox
from bokeh.palettes import Category20
from bokeh.plotting import curdoc
from holoviews.ipython.display_hooks import display
import colorcet as cc

colors = ['lat', 'lon', 'DOY', 'site', 'igbp']

def create_figure(color='lat', **kwargs):
    opts = {'plot': {'color_index': color, 'show_legend': False,
                     'width': 600, 'height': 600, 'colorbar': True,
                     'tools': ['hover']},
            'style': {'cmap': 'magma', 'legend': False}
}
    if color == 'DOY':
        opts['style']['cmap'] = cc.cm['cyclic_mrybm_35_75_c68']
    if color == 'igbp':
        opts['style']['cmap'] = 'Category20'
        opts['plot']['legend_position'] ='right'
        opts['plot']['show_legend'] = True
    if color == 'site':
        opts['style']['cmap'] = 'Category20'
        opts['plot']['colorbar'] = False
        opts['plot']['width'] = 700

    opts.update(**kwargs)
    chart = hv.Scatter(
        s, kdims=['x0', 'x1'], vdims=[color, 'site'], extents=(-15,-15,15,15)
    ).opts(plot=opts['plot'], style=opts['style'])
    return display(chart)

from ipywidgets import interactive

w = interactive(create_figure, color=colors)
w



## Taking a closer look at vegetation

In [20]:
igbp_vegetation = {
    'ENF': '01 - Evergreen Needleleaf forest',
    'EBF': '02 - Evergreen Broadleaf forest',
    'DNF': '03 - Deciduous Needleleaf forest',
    'DBF': '04 - Deciduous Broadleaf forest',
    'MF' : '05 - Mixed forest',
    'CSH': '06 - Closed shrublands',
    'OSH': '07 - Open shrublands',
    'WSA': '08 - Woody savannas',
    'SAV': '09 - Savannas',
    'GRA': '10 - Grasslands',
    'WET': '11 - Permanent wetlands',
    'CRO': '12 - Croplands',
}

In [21]:
s['vegetation'] = s['igbp'].apply(lambda x: igbp_vegetation[x])

In [22]:
ds = hv.Dataset(s, ['x0', 'vegetation'], ['x1', 'site'])
grouped = ds.to(hv.Scatter, kdims=['x0', 'x1'], extents=(-15,-15,15,15), vdims=['site'])

In [23]:
# https://lpdaac.usgs.gov/about/news_archive/modisterra_land_cover_types_yearly_l3_global_005deg_cmg_mod12c1
lpdaac_palette = [
    '#008000', '#00FF00', '#99CC00', '#99FF99', '#339966', '#993366',
    '#FFCC99', '#CCFFCC', '#FFCC00', '#FF9900', '#006699', '#FFFF00'
]

In [24]:
%%opts Scatter [width=800, height=600] (color=Cycle(lpdaac_palette), size=1, muted_alpha=0)
grouped.overlay('vegetation').options(legend_position='right')

Isolate each vegetation type so that any site eccentricities are made clear. In this, let's **color by site ID**

In [25]:
grouped.options(color_index='site', cmap='Category20', show_legend=False, size=1, alpha=0.8).layout().cols(3)

## Prediction
### One site

In [15]:
from sklearn.neighbors import NearestNeighbors
from sklearn.model_selection import LeaveOneGroupOut
from dateutil import rrule
from datetime import datetime, timedelta
from sklearn.preprocessing import StandardScaler
from sklearn.linear_model import LinearRegression

def fit_and_predict(X_train, y_train, X_test, nbrs=False):
    if nbrs:
        _nbrs = NearestNeighbors(n_neighbors=3, algorithm='ball_tree').fit(X_train)
        distances, indices = _nbrs.kneighbors(X_test)
    else:
        indices = np.arange(len(X_train), dtype=int)
    
    X_train_filtered = X_train[indices.flat[:]] 
    y_train_filtered = y_train[indices.flat[:]] 
        
    
    model = LinearRegression()
    model.fit(X_train_filtered, y_train_filtered)
    return model.predict(X_test)

In [16]:
def get_train_test_data(site, df, sites, doy):
    i = sites == site
    assert i.sum() > 0
    _df = df[i]
    year = (doy // 365)

    y = _df['y'].values
    X = pd.DataFrame({col: _df[col].values 
                      for col in show.columns 
                      if col != 'y'})
    return X.values[:-365], X.values[-365:], y[:-365], y[-365:]

data = []
for i, site in enumerate(df.site.unique()):
    if i % 20 == 0:
        print(i)
    X_train, X_test, y_train, y_test = get_train_test_data(site, show, sites=df.site, doy=df.DOY)
    y_hat = fit_and_predict(X_train, y_train, X_test)
    data += [{'site': site, 'corrcoef': np.corrcoef(y_hat, y_test)[0, 1]}]

0


  c /= stddev[:, None]
  c /= stddev[None, :]


20
40
60
80
100
120
140
160


In [46]:
stats = pd.DataFrame(data) 
stats.head()

corrs = stats.corrcoef[~np.isnan(stats.corrcoef.values)]
frequencies, edges = np.histogram(corrs, 20)

s1 = hv.Histogram((frequencies, edges), extents=(-1, None, 1, None), label='one site')
print("Correlation coefficients: mean={:0.3f}, median={:0.3}".format(np.mean(corrs), np.median(corrs)))
s1

Correlation coefficients: mean=0.497, median=0.573


### Multiple sites
Linear models work well *at one site* but this is confounded by

* lat/lon
* day of year
* environment type

In [49]:
assert 'site' not in show.columns
y = show['y'].values
X = pd.DataFrame({col: show[col].values 
                  for col in show.columns 
                  if col != 'y'})
print("X.shape =", X.shape)
assert 'y' not in X.columns

# transform data matrix so 0 mean, unit variance for each feature
X = StandardScaler().fit_transform(X.values)

X.shape = (532544, 28)


In [50]:
def prediction_stats(train_idx, test_idx, X, y, doy=None, predict_each='season', nbrs=False):
    start = datetime(2000, 1, 1)
    end = start + timedelta(days=365)
    
    if predict_each == 'month':
        get_time_id = lambda dt: dt.month
    elif predict_each == 'year':
        get_time_id = lambda dt: 1
    elif predict_each == 'season':
        seasons = {'spring': [3, 4, 5],
                   'summer': [6, 7, 8],
                   'fall': [9, 10, 11],
                   'winter': [12, 1, 2]}
        seasons = {month: season_id
                   for season_id, months in enumerate(seasons.values())
                   for month in months}
        get_time_id = lambda dt: seasons[dt.month] 
    else:
        msg = "predict_each should be in {'year', 'month', 'season'}, got '{}'"
        raise ValueError(msg.format(predict_each))
    
    # from https://stackoverflow.com/questions/153584/how-to-iterate-over-a-timespan-after-days-hours-weeks-and-months-in-python
    time_partitions = {(dt - start).days: get_time_id(dt)
                       for dt in rrule.rrule(rrule.DAILY, dtstart=start, until=end)}
    time_partitions[366] = max(time_partitions.values())
    
    test_days = doy[test_idx]
    
    preds = []
    for time_partition in time_partitions.values():
        if len(time_partitions.values()) > 1:
            time_idx = [i for i, day in enumerate(doy) if time_partitions[day] == time_partition]
            
            # get the test set specific to this time instance
            time_test_idx = np.intersect1d(test_idx, time_idx)
        else:
            time_test_idx = test_idx 

        if len(time_test_idx) == 0:
            continue
            
        y_hat = fit_and_predict(X[train_idx], y[train_idx], X[time_test_idx], nbrs=nbrs)
        y_test = y[time_test_idx]
        preds += [{'predicted': y_hat,
                   'actual': y_test,
                   'time_partition': time_partition,
                   'corrcoef': np.corrcoef(y_hat, y_test)[0][1]}]
    actual = [p['actual'] for p in preds]
    predicted = [p['predicted'] for p in preds]
    actual = np.concatenate(actual).flat[:]
    predicted = np.concatenate(predicted).flat[:]
    return {'time_partitions': preds,
            'actual': actual,
            'predicted': predicted,
            'corrcoef': np.corrcoef(actual, predicted)[0][1]}


from sklearn.model_selection import LeaveOneGroupOut
sep = LeaveOneGroupOut()
# train_idx, test_idx = list(sep.split(X, y, sites))[0]
# _ = prediction_stats(train_idx, test_idx, X, y, doy=dropped['DOY'])

In [51]:
from sklearn.model_selection import LeaveOneGroupOut
sep = LeaveOneGroupOut()
corrs = []

futures = []
n_splits = sep.get_n_splits(X, y, sites)
X_future = client.scatter(X)
y_future = client.scatter(y)
doy_future = client.scatter(dropped['DOY'])
for i, (train_index, test_index) in enumerate(sep.split(X, y, sites)):
    futures += [{'site_id': i,
                 'train_index': train_index,
                 'test_index': test_index,
                 'stats': client.submit(prediction_stats,
                                        train_index,
                                        test_index,
                                        X_future,
                                        y_future,
                                        doy=doy_future)}]

In [None]:
results = client.gather(futures)

In [83]:
out = [{'site_id': result['site_id'], **result['stats']}
       for result in results] 

In [19]:
# df = pd.DataFrame(out)
# df[['corrcoef', 'site_id']].to_csv('2018-09-30-global-exps-corrcoef.csv')
df = pd.read_csv('2018-09-30-global-exps-corrcoef.csv')
df.head()

Unnamed: 0.1,Unnamed: 0,corrcoef,site_id
0,0,0.314836,0
1,1,0.214287,1
2,2,0.309658,2
3,3,0.022033,3
4,4,0.14236,4


In [44]:
%%opts VLine [show_legend=False] VLine (color='red')
corrs = df.corrcoef[~np.isnan(df.corrcoef.values)]
frequencies, edges = np.histogram(corrs, 20)

print("Correlation coefficients: mean={:0.3f}, median={:0.3}".format(np.mean(corrs), np.median(corrs)))
c1 = hv.Histogram((frequencies, edges), extents=(-1, None, 1, None), label='global')
c1

Correlation coefficients: mean=0.398, median=0.441


In [49]:
%%opts Histogram (alpha=0.5) Overlay [legend_position='top_left']
c1 * s1