In [1]:
import numpy as np
import pandas as pd
import scipy
from cuml import KalmanFilter
import cudf
import os

# Helper Functions

In [2]:
from timeit import default_timer

class Timer(object):
    def __init__(self):
        self._timer = default_timer
    
    def __enter__(self):
        self.start()
        return self

    def __exit__(self, *args):
        self.stop()

    def start(self):
        """Start the timer."""
        self.start = self._timer()

    def stop(self):
        """Stop the timer. Calculate the interval in seconds."""
        self.end = self._timer()
        self.interval = self.end - self.start

In [3]:
import gzip
def load_data(nrows, ncols, cached = '/rapids/notebooks/cumldata/mortgage.npy.gz',source='mortgage'):
    if os.path.exists(cached) and source=='mortgage':
        print('use mortgage data')
        with gzip.open(cached) as f:
            X = np.load(f)
        X = X[np.random.randint(0,X.shape[0]-1,nrows),:ncols]
    else:
        print('use random data')
        X = np.random.random((nrows,ncols)).astype('float32')
    df = pd.DataFrame({'fea%d'%i:X[:,i] for i in range(X.shape[1])}).fillna(0)
    return df

In [4]:
from sklearn.metrics import mean_squared_error
def array_equal(a,b,threshold=1e-2,with_sign=True,metric='mse'):
    a = to_nparray(a)
    b = to_nparray(b)
    if with_sign == False:
        a,b = np.abs(a),np.abs(b)
    if metric=='mse':
        error = mean_squared_error(a,b)
    else:
        error = np.sum(a!=b)/(a.shape[0]*a.shape[1])
    res = error<threshold
    return res

def to_nparray(x):
    if isinstance(x,np.ndarray) or isinstance(x,pd.DataFrame):
        return np.array(x)
    elif isinstance(x,np.float64):
        return np.array([x])
    elif isinstance(x,cudf.DataFrame) or isinstance(x,cudf.Series):
        return x.to_pandas().values
    return x    

In [5]:
def spKalman(data, x, z, n_iter = 50):
    # intial parameters
    sz = data # size of array
    x = -0.37727 # truth value (typo in example at top of p. 13 calls this z)
    z = np.random.normal(x,0.1,size=sz) # observations (normal about x, sigma=0.1)

    Q = 1e-5 # process variance

    # allocate space for arrays
    xhat=np.zeros(sz)      # a posteri estimate of x
    P=np.zeros(sz)         # a posteri error estimate
    xhatminus=np.zeros(sz) # a priori estimate of x
    Pminus=np.zeros(sz)    # a priori error estimate
    K=np.zeros(sz)         # gain or blending factor

    R = 0.1**2 # estimate of measurement variance, change to see effect

    # intial guesses
    xhat[0] = 0.0
    P[0] = 1.0

    for k in range(1,n_iter):
        # time update
        xhatminus[k] = xhat[k-1]
        Pminus[k] = P[k-1]+Q

        # measurement update
        K[k] = Pminus[k]/( Pminus[k]+R )
        xhat[k] = xhatminus[k]+K[k]*(z[k]-xhatminus[k])
        P[k] = (1-K[k])*Pminus[k]
    return Pminus 

In [7]:
def cuMLKalman(data, dim_x,dim_z, n_iter = 50):
    f = KalmanFilter(dim_x, dim_z)
    f.x = np.array(dim_x, 1)   # velocity
    f.F = np.array([[1.,1.], [0.,1.]])
    f.H = np.array(dim_z, dim_x)
    f.P = np.array(dim_x, dim_z)
    f.R = 5

    for k in range(1,n_iter):
        z = numba.cuda.to_device(np.array([i]))
        f.predict()
        f.update(z)

# Run tests

In [None]:
%%time
nrows = 2**16
ncols = 40
n_iter = 50 

X = load_data(nrows,ncols)
#spKalman(X, nrows, ncols, n_iter)
cuMLKalman(X, nrows, ncols, n_iter)

print('data',X.shape)

In [None]:
n_neighbors = 10

In [None]:
%%time
knn_sk = skKNN(X)
D_sk,I_sk = knn_sk.query(X,n_neighbors)

In [None]:
%%time
X = cudf.DataFrame.from_pandas(X)

In [None]:
%%time
knn_cuml = cumlKNN(n_gpus=1)
knn_cuml.fit(X)
D_cuml,I_cuml = knn_cuml.query(X,n_neighbors)

In [None]:
passed = array_equal(D_sk,D_cuml)
message = 'compare knn: cuml vs sklearn distances %s'%('equal'if passed else 'NOT equal')
print(message)
passed = array_equal(I_sk,I_cuml)
message = 'compare knn: cuml vs sklearn indexes %s'%('equal'if passed else 'NOT equal')
print(message)