diff --git a/requirements.txt b/requirements.txt index 76e37d6a..86e57035 100644 --- a/requirements.txt +++ b/requirements.txt @@ -19,6 +19,7 @@ certifi==2019.3.9 cffi==1.12.2 chardet==3.0.4 Click==7.0 +cloudpickle==1.2.2 colorama==0.4.1 colorclass==2.2.0 cookies==2.2.1 @@ -27,7 +28,7 @@ coreschema==0.0.4 coverage==4.5.3 cryptography==2.6.1 cycler==0.10.0 -dask==1.1.4 +dask==2.6.0 decorator==4.4.0 defusedxml==0.5.0 Django==2.1.7 @@ -39,7 +40,7 @@ django-cookie-law==2.0.1 django-cors-headers==2.5.2 django-countries==5.3.3 django-datetime-widget==0.9.3 -django-debug-toolbar==1.11 +django-debug-toolbar==2.0 django-extensions==2.1.6 django-htmlmin==0.11.0 django-leaflet==0.24.0 @@ -64,8 +65,10 @@ entrypoints==0.3 execnet==1.5.0 factory-boy==2.11.1 Faker==1.0.4 +fastparquet==0.3.2 fitparse==1.1.0 Flask==1.0.2 +fsspec==0.5.2 future==0.17.1 geocoder==1.38.1 geos==0.2.1 @@ -74,6 +77,7 @@ html5lib==1.0.1 htmlmin==0.1.12 HTMLParser==0.0.2 httplib2==0.12.1 +hvplot==0.4.0 icalendar==4.0.3 idna==2.8 image==1.5.27 @@ -99,10 +103,12 @@ jupyterlab-server==0.3.0 keyring==18.0.0 kiwisolver==1.0.1 kombu==4.5.0 +llvmlite==0.30.0 lxml==4.3.2 Markdown==3.0.1 MarkupSafe==1.1.1 matplotlib==3.0.3 +minify==0.1.4 MiniMockTest==0.5 mistune==0.8.4 mock==2.0.0 @@ -111,9 +117,11 @@ mpld3==0.3 mysqlclient==1.4.2.post1 nbconvert==5.4.1 nbformat==4.4.0 +newrelic==5.2.1.129 nose==1.3.7 nose-parameterized==0.6.0 notebook==5.7.6 +numba==0.46.0 numpy==1.16.2 oauth2==1.9.0.post1 oauthlib==3.0.1 @@ -135,6 +143,7 @@ prompt-toolkit==2.0.9 psycopg2==2.8.1 ptyprocess==0.6.0 py==1.8.0 +pyarrow==0.15.0 pycparser==2.19 Pygments==2.3.1 pyparsing==2.3.1 @@ -160,7 +169,7 @@ ratelim==0.1.6 redis==3.2.1 requests==2.21.0 requests-oauthlib==1.2.0 -rowingdata==2.5.4 +rowingdata==2.5.5 rowingphysics==0.5.0 rq==0.13.0 scipy==1.2.1 @@ -179,7 +188,9 @@ terminado==0.8.1 terminaltables==3.1.0 testpath==0.4.2 text-unidecode==1.2 +thrift==0.11.0 timezonefinder==4.0.1 +toolz==0.10.0 tornado==6.0.1 tqdm==4.31.1 traitlets==4.3.2 @@ -196,3 +207,4 @@ xlrd==1.2.0 xmltodict==0.12.0 yamjam==0.1.7 yamllint==1.15.0 +yuicompressor==2.4.8 diff --git a/rowers/c2stuff.py b/rowers/c2stuff.py index 4d034e6a..5c71d85a 100644 --- a/rowers/c2stuff.py +++ b/rowers/c2stuff.py @@ -29,6 +29,70 @@ queue = django_rq.get_queue('default') queuelow = django_rq.get_queue('low') queuehigh = django_rq.get_queue('low') from rowers.utils import myqueue +from rowers.models import C2WorldClassAgePerformance + +def getagegrouprecord(age,sex='male',weightcategory='hwt', + distance=2000,duration=None,indf=pd.DataFrame()): + + if not indf.empty: + if not duration: + df = indf[indf['distance'] == distance] + else: + duration = 60*int(duration) + df = indf[indf['duration'] == duration] + else: + if not duration: + df = pd.DataFrame( + list( + C2WorldClassAgePerformance.objects.filter( + distance=distance, + sex=sex, + weightcategory=weightcategory + ).values() + ) + ) + else: + duration=60*int(duration) + df = pd.DataFrame( + list( + C2WorldClassAgePerformance.objects.filter( + duration=duration, + sex=sex, + weightcategory=weightcategory + ).values() + ) + ) + + if not df.empty: + ages = df['age'] + powers = df['power'] + + #poly_coefficients = np.polyfit(ages,powers,6) + fitfunc = lambda pars, x: np.abs(pars[0])*(1-x/max(120,pars[1]))-np.abs(pars[2])*np.exp(-x/np.abs(pars[3]))+np.abs(pars[4])*(np.sin(np.pi*x/max(50,pars[5]))) + errfunc = lambda pars, x,y: fitfunc(pars,x)-y + + p0 = [700,120,700,10,100,100] + + try: + p1, success = optimize.leastsq(errfunc,p0[:], + args = (ages,powers)) + except: + p1 = p0 + success = 0 + + if success: + power = fitfunc(p1, float(age)) + + #power = np.polyval(poly_coefficients,age) + + power = 0.5*(np.abs(power)+power) + else: + power = 0 + else: + power = 0 + + return power + oauth_data = { 'client_id': C2_CLIENT_ID, diff --git a/rowers/dataprep.py b/rowers/dataprep.py index 9d43a330..1fb142d0 100644 --- a/rowers/dataprep.py +++ b/rowers/dataprep.py @@ -4,11 +4,10 @@ from __future__ import print_function from __future__ import unicode_literals - # All the data preparation, data cleaning and data mangling should # be defined here from __future__ import unicode_literals, absolute_import -from rowers.models import Workout, StrokeData,Team +from rowers.models import Workout, Team import pytz @@ -16,6 +15,7 @@ from rowingdata import rowingdata as rrdata from rowingdata import rower as rrower +import shutil from shutil import copyfile from rowingdata import ( @@ -26,6 +26,10 @@ from rowers.tasks import handle_sendemail_unrecognized from rowers.tasks import handle_zip_file from pandas import DataFrame, Series +import dask.dataframe as dd +from dask.delayed import delayed +import pyarrow.parquet as pq +import pyarrow as pa from django.utils import timezone from django.utils.timezone import get_current_timezone @@ -51,7 +55,7 @@ from rowingdata import ( from rowingdata.csvparsers import HumonParser -from rowers.metrics import axes,calc_trimp,rowingmetrics +from rowers.metrics import axes,calc_trimp,rowingmetrics,dtypes from rowers.models import strokedatafields #allowedcolumns = [item[0] for item in rowingmetrics] @@ -113,6 +117,7 @@ columndict = { 'cumdist': 'cum_dist', } + from scipy.signal import savgol_filter import datetime @@ -348,22 +353,23 @@ def clean_df_stats(datadf, workstrokesonly=True, ignorehr=True, ignoreadvanced=False): # clean data remove zeros and negative values + # bring metrics which have negative values to positive domain - if datadf.empty: + if len(datadf)==0: return datadf try: datadf['catch'] = -datadf['catch'] - except KeyError: + except (KeyError,TypeError): pass try: datadf['peakforceangle'] = datadf['peakforceangle'] + 1000 - except KeyError: + except (KeyError,TypeError): pass try: datadf['hr'] = datadf['hr'] + 10 - except KeyError: + except (KeyError,TypeError): pass # protect 0 spm values from being nulled @@ -378,6 +384,7 @@ def clean_df_stats(datadf, workstrokesonly=True, ignorehr=True, pass datadf.replace(to_replace=0, value=np.nan, inplace=True) + # datadf = datadf.map_partitions(lambda df:df.replace(to_replace=0,value=np.nan)) # bring spm back to real values try: @@ -388,141 +395,141 @@ def clean_df_stats(datadf, workstrokesonly=True, ignorehr=True, # return from positive domain to negative try: datadf['catch'] = -datadf['catch'] - except KeyError: + except (KeyError,TypeError): pass try: datadf['peakforceangle'] = datadf['peakforceangle'] - 1000 - except KeyError: + except (KeyError,TypeError): pass try: datadf['hr'] = datadf['hr'] - 10 - except KeyError: + except (KeyError,TypeError): pass # clean data for useful ranges per column if not ignorehr: try: mask = datadf['hr'] < 30 - datadf.loc[mask, 'hr'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['spm'] < 0 - datadf.loc[mask,'spm'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['efficiency'] > 200. - datadf.loc[mask, 'efficiency'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['spm'] < 10 - datadf.loc[mask, 'spm'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['pace'] / 1000. > 300. - datadf.loc[mask, 'pace'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['efficiency'] < 0. - datadf.loc[mask, 'efficiency'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['pace'] / 1000. < 60. - datadf.loc[mask, 'pace'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['spm'] > 60 - datadf.loc[mask, 'spm'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: - mask = datadf['wash'] < 1 + mask = datadf['wash'] > 1 datadf.loc[mask, 'wash'] = np.nan - except KeyError: + except (KeyError,TypeError): pass if not ignoreadvanced: try: mask = datadf['rhythm'] < 5 - datadf.loc[mask, 'rhythm'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['rhythm'] > 70 - datadf.loc[mask, 'rhythm'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['power'] < 20 - datadf.loc[mask, 'power'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['drivelength'] < 0.5 - datadf.loc[mask, 'drivelength'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['forceratio'] < 0.2 - datadf.loc[mask, 'forceratio'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['forceratio'] > 1.0 - datadf.loc[mask, 'forceratio'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['drivespeed'] < 0.5 - datadf.loc[mask, 'drivespeed'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['drivespeed'] > 4 - datadf.loc[mask, 'drivespeed'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['driveenergy'] > 2000 - datadf.loc[mask, 'driveenergy'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['driveenergy'] < 100 - datadf.loc[mask, 'driveenergy'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass try: mask = datadf['catch'] > -30. - datadf.loc[mask, 'catch'] = np.nan - except KeyError: + datadf.mask(mask,inplace=True) + except (KeyError,TypeError): pass workoutstateswork = [1, 4, 5, 8, 9, 6, 7] @@ -569,41 +576,6 @@ def getstatsfields(): return fieldlist, fielddict -def getstatsfields_old(): - # Get field names and remove those that are not useful in stats - fields = StrokeData._meta.get_fields() - - fielddict = {field.name: field.verbose_name for field in fields} - - # fielddict.pop('workoutid') - fielddict.pop('ergpace') - fielddict.pop('hr_an') - fielddict.pop('hr_tr') - fielddict.pop('hr_at') - fielddict.pop('hr_ut2') - fielddict.pop('hr_ut1') - fielddict.pop('time') - fielddict.pop('distance') - fielddict.pop('nowindpace') - fielddict.pop('fnowindpace') - fielddict.pop('fergpace') - fielddict.pop('equivergpower') -# fielddict.pop('workoutstate') - fielddict.pop('fpace') - fielddict.pop('pace') - fielddict.pop('id') - fielddict.pop('ftime') - fielddict.pop('x_right') - fielddict.pop('hr_max') - fielddict.pop('hr_bottom') - fielddict.pop('cumdist') - - try: - fieldlist = [field for field, value in fielddict.iteritems()] - except AttributeError: - fieldlist = [field for field, value in fielddict.items()] - - return fieldlist, fielddict # A string representation for time deltas @@ -1616,75 +1588,7 @@ def new_workout_from_df(r, df, return (id, message) -# Compare the data from the CSV file and the database -# Currently only calculates number of strokes. To be expanded with -# more elaborate testing if needed -def compare_data(id): - row = Workout.objects.get(id=id) - f1 = row.csvfilename - try: - rowdata = rdata(f1) - l1 = len(rowdata.df) - except AttributeError: - rowdata = 0 - l1 = 0 - engine = create_engine(database_url, echo=False) - query = sa.text('SELECT COUNT(*) FROM strokedata WHERE workoutid={id};'.format( - id=id, - )) - with engine.connect() as conn, conn.begin(): - try: - res = conn.execute(query) - l2 = res.fetchall()[0][0] - except: - print("Database Locked") - conn.close() - engine.dispose() - lfile = l1 - ldb = l2 - return l1 == l2 and l1 != 0, ldb, lfile - -# Repair data for workouts where the CSV file is lost (or the DB entries -# don't exist) - - -def repair_data(verbose=False): - ws = Workout.objects.all() - for w in ws: - if verbose: - sys.stdout.write(".") - test, ldb, lfile = compare_data(w.id) - if not test: - if verbose: - print(w.id, lfile, ldb) - try: - rowdata = rdata(w.csvfilename) - if rowdata and len(rowdata.df): - update_strokedata(w.id, rowdata.df) - - except (IOError, AttributeError): - pass - - if lfile == 0: - # if not ldb - delete workout - - try: - data = read_df_sql(w.id) - try: - datalength = len(data) - except AttributeError: - datalength = 0 - - if datalength != 0: - data.rename(columns=columndict, inplace=True) - res = data.to_csv(w.csvfilename + '.gz', - index_label='index', - compression='gzip') - else: - w.delete() - except: - pass # A wrapper around the rowingdata class, with some error catching @@ -1710,17 +1614,11 @@ def rdata(file, rower=rrower()): def delete_strokedata(id): - engine = create_engine(database_url, echo=False) - query = sa.text('DELETE FROM strokedata WHERE workoutid={id};'.format( - id=id, - )) - with engine.connect() as conn, conn.begin(): - try: - result = conn.execute(query) - except: - print("Database Locked") - conn.close() - engine.dispose() + dirname = 'media/strokedata_{id}.parquet.gz'.format(id=id) + try: + shutil.rmtree(dirname) + except FileNotFoundError: + pass # Replace stroke data in DB with data from CSV file @@ -1747,7 +1645,6 @@ def testdata(time, distance, pace, spm): def getrowdata_db(id=0, doclean=False, convertnewtons=True, checkefficiency=True): data = read_df_sql(id) - data['x_right'] = data['x_right'] / 1.0e6 data['deltat'] = data['time'].diff() if data.empty: @@ -1771,8 +1668,111 @@ def getrowdata_db(id=0, doclean=False, convertnewtons=True, # Fetch a subset of the data from the DB +def getsmallrowdata_db(columns, ids=[], doclean=True,workstrokesonly=True,compute=True): + # prepmultipledata(ids) -def getsmallrowdata_db(columns, ids=[], doclean=True, workstrokesonly=True): + csvfilenames = ['media/strokedata_{id}.parquet.gz'.format(id=id) for id in ids] + data = [] + columns = [c for c in columns if c != 'None'] + columns = list(set(columns)) + + if len(ids)>1: + for id,f in zip(ids,csvfilenames): + try: + #df = dd.read_parquet(f,columns=columns,engine='pyarrow') + df = pd.read_parquet(f,columns=columns) + data.append(df) + except OSError: + rowdata, row = getrowdata(id=id) + if rowdata and len(rowdata.df): + datadf = dataprep(rowdata.df,id=id,bands=True,otwpower=True,barchart=True) + # df = dd.read_parquet(f,columns=columns,engine='pyarrow') + df = pd.read_parquet(f,columns=columns) + data.append(df) + + df = pd.concat(data,axis=0) + # df = dd.concat(data,axis=0) + + else: + try: + df = pd.read_parquet(csvfilenames[0],columns=columns) + except OSError: + rowdata,row = getrowdata(id=ids[0]) + if rowdata and len(rowdata.df): + data = dataprep(rowdata.df,id=ids[0],bands=True,otwpower=True,barchart=True) + df = pd.read_parquet(csvfilenames[0],columns=columns) + # df = dd.read_parquet(csvfilenames[0], + # column=columns,engine='pyarrow', + # ) + + # df = df.loc[:,~df.columns.duplicated()] + + + + if compute: + data = df.copy() + if doclean: + data = clean_df_stats(data, ignorehr=True, + workstrokesonly=workstrokesonly) + data.dropna(axis=1,how='all',inplace=True) + data.dropna(axis=0,how='any',inplace=True) + return data + + return df + +def getsmallrowdata_db_dask(columns, ids=[], doclean=True,workstrokesonly=True,compute=True): + # prepmultipledata(ids) + + csvfilenames = ['media/strokedata_{id}.parquet.gz'.format(id=id) for id in ids] + data = [] + columns = [c for c in columns if c != 'None'] + columns = list(set(columns)) + + if len(ids)>1: + for id,f in zip(ids,csvfilenames): + try: + #df = dd.read_parquet(f,columns=columns,engine='pyarrow') + df = dd.read_parquet(f,columns=columns) + data.append(df) + except OSError: + rowdata, row = getrowdata(id=id) + if rowdata and len(rowdata.df): + datadf = dataprep(rowdata.df,id=id,bands=True,otwpower=True,barchart=True) + # df = dd.read_parquet(f,columns=columns,engine='pyarrow') + df = dd.read_parquet(f,columns=columns) + data.append(df) + + df = dd.concat(data,axis=0) + # df = dd.concat(data,axis=0) + + else: + try: + df = dd.read_parquet(csvfilenames[0],columns=columns) + except OSError: + rowdata,row = getrowdata(id=ids[0]) + if rowdata and len(rowdata.df): + data = dataprep(rowdata.df,id=ids[0],bands=True,otwpower=True,barchart=True) + df = dd.read_parquet(csvfilenames[0],columns=columns) + # df = dd.read_parquet(csvfilenames[0], + # column=columns,engine='pyarrow', + # ) + + # df = df.loc[:,~df.columns.duplicated()] + + + + if compute: + data = df.compute() + if doclean: + data = clean_df_stats(data, ignorehr=True, + workstrokesonly=workstrokesonly) + data.dropna(axis=1,how='all',inplace=True) + data.dropna(axis=0,how='any',inplace=True) + return data + + return df + +def getsmallrowdata_db_old(columns, ids=[], doclean=True, workstrokesonly=True): prepmultipledata(ids) data,extracols = read_cols_df_sql(ids, columns) if extracols and len(ids)==1: @@ -1850,31 +1850,20 @@ def getrowdata(id=0): # safety net for programming errors elsewhere in the app # Also used heavily when I moved from CSV file only to CSV+Stroke data +import glob def prepmultipledata(ids, verbose=False): - query = sa.text('SELECT DISTINCT workoutid FROM strokedata') - engine = create_engine(database_url, echo=False) + filenames = glob.glob('media/*.parquet') + ids = [id for id in ids if 'media/strokedata_{id}.parquet.gz'.format(id=id) not in filenames] - with engine.connect() as conn, conn.begin(): - res = conn.execute(query) - res = list(itertools.chain.from_iterable(res.fetchall())) - conn.close() - engine.dispose() - - try: - ids2 = [int(id) for id in ids] - except ValueError: - ids2 = ids - - res = list(set(ids2) - set(res)) - for id in res: + for id in ids: rowdata, row = getrowdata(id=id) if verbose: print(id) if rowdata and len(rowdata.df): data = dataprep(rowdata.df, id=id, bands=True, barchart=True, otwpower=True) - return res + return ids # Read a set of columns for a set of workout ids, returns data as a # pandas dataframe @@ -1883,6 +1872,66 @@ def prepmultipledata(ids, verbose=False): def read_cols_df_sql(ids, columns, convertnewtons=True): # drop columns that are not in offical list # axx = [ax[0] for ax in axes] + + extracols = [] + + columns = list(columns) + ['distance', 'spm', 'workoutid'] + columns = [x for x in columns if x != 'None'] + columns = list(set(columns)) + ids = [int(id) for id in ids] + + + if len(ids) == 0: + return pd.DataFrame(),extracols + elif len(ids) == 1: + try: + filename = 'media/strokedata_{id}.parquet.gz'.format(id=ids[0]) + df = pd.read_parquet(filename,columns=columns) + except OSError: + rowdata,row = getrowdata(id=ids[0]) + if rowdata and len(rowdata.df): + datadf = dataprep(rowdata.df,id=ids[0],bands=True,otwpower=True,barchart=True) + df = pd.read_parquet(filename,columns=columns) + else: + data = [] + filenames = ['media/strokedata_{id}.parquet.gz'.format(id=id) for id in ids] + for id,f in zip(ids,filenames): + try: + df = pd.read_parquet(f,columns=columns) + data.append(df) + except OSError: + rowdata,row = getrowdata(id=id) + if rowdata and len(rowdata.df): + datadf = dataprep(rowdata.df,id=id,bands=True,otwpower=True,barchart=True) + df = pd.read_parquet(f,columns=columns) + data.append(df) + + df = pd.concat(data,axis=0) + + + df = df.fillna(value=0) + + if 'peakforce' in columns: + funits = ((w.id, w.forceunit) + for w in Workout.objects.filter(id__in=ids)) + for id, u in funits: + if u == 'lbs': + mask = df['workoutid'] == id + df.loc[mask, 'peakforce'] = df.loc[mask, 'peakforce'] * lbstoN + if 'averageforce' in columns: + funits = ((w.id, w.forceunit) + for w in Workout.objects.filter(id__in=ids)) + for id, u in funits: + if u == 'lbs': + mask = df['workoutid'] == id + df.loc[mask, 'averageforce'] = df.loc[mask, + 'averageforce'] * lbstoN + + return df,extracols + +def read_cols_df_sql_old(ids, columns, convertnewtons=True): + # drop columns that are not in offical list + # axx = [ax[0] for ax in axes] prepmultipledata(ids) axx = [f.name for f in StrokeData._meta.get_fields()] @@ -1949,8 +1998,34 @@ def read_cols_df_sql(ids, columns, convertnewtons=True): # Read stroke data from the DB for a Workout ID. Returns a pandas dataframe - def read_df_sql(id): + try: + f = 'media/strokedata_{id}.parquet.gz'.format(id=id) + df = pd.read_parquet(f) + except OSError: + rowdata,row = getrowdata(id=ids[0]) + if rowdata and len(rowdata.df): + data = dataprep(rowdata.df,id=ids[0],bands=True,otwpower=True,barchart=True) + df = pd.read_parquet(f) + + df = df.fillna(value=0) + + funit = Workout.objects.get(id=id).forceunit + + if funit == 'lbs': + try: + df['peakforce'] = df['peakforce'] * lbstoN + except KeyError: + pass + + try: + df['averageforce'] = df['averageforce'] * lbstoN + except KeyError: + pass + + return df + +def read_df_sql_old(id): engine = create_engine(database_url, echo=False) df = pd.read_sql_query(sa.text('SELECT * FROM strokedata WHERE workoutid={id} ORDER BY time ASC'.format( @@ -2142,14 +2217,13 @@ def add_efficiency(id=0): rowdata = rowdata.fillna(method='ffill') delete_strokedata(id) + if id != 0: rowdata['workoutid'] = id - engine = create_engine(database_url, echo=False) - with engine.connect() as conn, conn.begin(): - rowdata.to_sql('strokedata', engine, - if_exists='append', index=False) - conn.close() - engine.dispose() + filename = 'media/strokedata_{id}.parquet.gz'.format(id=id) + df = dd.from_pandas(rowdata,npartitions=1) + df.to_parquet(filename,engine='fastparquet',compression='GZIP') + return rowdata # This is the main routine. @@ -2292,19 +2366,6 @@ def dataprep(rowdatadf, id=0, bands=True, barchart=True, otwpower=True, except KeyError: rowdatadf[' ElapsedTime (sec)'] = rowdatadf['TimeStamp (sec)'] - if barchart: - # time increments for bar chart - time_increments = rowdatadf.loc[:, ' ElapsedTime (sec)'].diff() - try: - time_increments.iloc[0] = time_increments.iloc[1] - except (KeyError, IndexError): - time_increments.iloc[0] = 1. - - time_increments = 0.5 * time_increments + 0.5 * np.abs(time_increments) - x_right = (t2 + time_increments.apply(lambda x: timedeltaconv(x))) - - data['x_right'] = x_right - if empower: try: wash = rowdatadf.loc[:, 'wash'] @@ -2441,12 +2502,17 @@ def dataprep(rowdatadf, id=0, bands=True, barchart=True, otwpower=True, # write data if id given if id != 0: data['workoutid'] = id + data.fillna(0,inplace=True) + data = data.astype( + dtype=dtypes, + ) + + + filename = 'media/strokedata_{id}.parquet.gz'.format(id=id) + df = dd.from_pandas(data,npartitions=1) + df.to_parquet(filename,engine='fastparquet',compression='GZIP') + - engine = create_engine(database_url, echo=False) - with engine.connect() as conn, conn.begin(): - data.to_sql('strokedata', engine, if_exists='append', index=False) - conn.close() - engine.dispose() return data diff --git a/rowers/dataprepnodjango.py b/rowers/dataprepnodjango.py index 9fd3d196..c7e2871c 100644 --- a/rowers/dataprepnodjango.py +++ b/rowers/dataprepnodjango.py @@ -16,6 +16,8 @@ from pandas import DataFrame,Series import pandas as pd import numpy as np import itertools +import dask.dataframe as dd +from dask.delayed import delayed from sqlalchemy import create_engine import sqlalchemy as sa @@ -145,6 +147,7 @@ def rdata(file,rower=rrower()): return res from rowers.utils import totaltime_sec_to_string +from rowers.metrics import dtypes # Creates C2 stroke data @@ -635,20 +638,11 @@ def new_workout_from_file(r,f2, return (id,message,f2) def delete_strokedata(id,debug=False): - if debug: - engine = create_engine(database_url_debug, echo=False) - else: - engine = create_engine(database_url, echo=False) - query = sa.text('DELETE FROM strokedata WHERE workoutid={id};'.format( - id=id, - )) - with engine.connect() as conn, conn.begin(): - try: - result = conn.execute(query) - except: - print("Database Locked") - conn.close() - engine.dispose() + dirname = 'media/strokedata_{id}.parquet.gz'.format(id=id) + try: + shutil.rmtree(dirname) + except FileNotFoundError: + pass def update_strokedata(id,df,debug=False): delete_strokedata(id,debug=debug) @@ -714,10 +708,25 @@ def testdata(time,distance,pace,spm): def getsmallrowdata_db(columns,ids=[],debug=False): + csvfilenames = ['media/strokedata_{id}.parquet'.format(id=id) for id in ids] + data = [] + columns = [c for c in columns if c != 'None'] - data = read_cols_df_sql(ids,columns,debug=debug) + if len(ids)>1: + for f in csvfilenames: + try: + df = pd.read_parquet(f,columns=columns,engine='pyarrow') + data.append(df) + except OSError: + pass - return data + + df = pd.concat(data,axis=0) + else: + df = pd.read_parquet(csvfilenames[0],columns=columns,engine='pyarrow') + + + return df def fitnessmetric_to_sql(m,table='powertimefitnessmetric',debug=False, doclean=False): @@ -761,51 +770,42 @@ def read_cols_df_sql(ids,columns,debug=False): columns = list(columns)+['distance','spm'] columns = [x for x in columns if x != 'None'] columns = list(set(columns)) - cls = '' + ids = [int(id) for id in ids] - if debug: - engine = create_engine(database_url_debug, echo=False) - else: - engine = create_engine(database_url, echo=False) - for column in columns: - cls += column+', ' - cls = cls[:-2] if len(ids) == 0: - query = sa.text('SELECT {columns} FROM strokedata WHERE workoutid=0'.format( - columns = cls, - )) + return pd.DataFrame() elif len(ids) == 1: - query = sa.text('SELECT {columns} FROM strokedata WHERE workoutid={id}'.format( - id = ids[0], - columns = cls, - )) + try: + filename = 'media/strokedata_{id}.parquet.gz'.format(id=ids[0]) + df = pd.read_parquet(filename,columns=columns) + except OSError: + pass else: - query = sa.text('SELECT {columns} FROM strokedata WHERE workoutid IN {ids}'.format( - columns = cls, - ids = tuple(ids), - )) - - df = pd.read_sql_query(query,engine) - engine.dispose() + data = [] + filenames = ['media/strokedata_{id}.parquet.gz'.format(id=id) for id in ids] + for id,f in zip(ids,filenames): + try: + df = pd.read_parquet(f,columns=columns) + data.append(df) + except OSError: + pass + + df = pd.concat(data,axis=0) + return df def read_df_sql(id,debug=False): - if debug: - engine = create_engine(database_url_debug, echo=False) - print("read_df",id) - print(database_url_debug) - else: - engine = create_engine(database_url, echo=False) + try: + f = 'media/strokedata_{id}.parquet.gz'.format(id=id) + df = pd.read_parquet(f) + except OSError: + pass - df = pd.read_sql_query(sa.text( - 'SELECT * FROM strokedata WHERE workoutid={id}'.format( - id=id - )), engine) + df = df.fillna(value=0) - engine.dispose() return df def getcpdata_sql(rower_id,table='cpdata',debug=False): @@ -1101,18 +1101,6 @@ def dataprep(rowdatadf,id=0,bands=True,barchart=True,otwpower=True, except KeyError: rowdatadf[' ElapsedTime (sec)'] = rowdatadf['TimeStamp (sec)'] - if barchart: - # time increments for bar chart - time_increments = rowdatadf.loc[:,' ElapsedTime (sec)'].diff() - try: - time_increments.iloc[0] = time_increments.iloc[1] - except (KeyError, IndexError): - time_increments.iloc[1] = 1. - - time_increments = 0.5*time_increments+0.5*np.abs(time_increments) - x_right = (t2+time_increments.apply(lambda x:timedeltaconv(x))) - - data['x_right'] = x_right if empower: try: @@ -1260,15 +1248,9 @@ def dataprep(rowdatadf,id=0,bands=True,barchart=True,otwpower=True, # write data if id given if id != 0: data['workoutid'] = id - - if debug: - engine = create_engine(database_url_debug, echo=False) - else: - engine = create_engine(database_url, echo=False) + data = data.astype(dtype=dtypes) + filename = 'media/strokedata_{id}.parquet.gz'.format(id=id) + df = dd.from_pandas(data,npartitions=1) + df.to_parquet(filename,engine='fastparquet',compression='GZIP') - with engine.connect() as conn, conn.begin(): - data.to_sql('strokedata',engine,if_exists='append',index=False) - - conn.close() - engine.dispose() return data diff --git a/rowers/interactiveplots.py b/rowers/interactiveplots.py index 21d48b40..95c5b5ed 100644 --- a/rowers/interactiveplots.py +++ b/rowers/interactiveplots.py @@ -77,6 +77,7 @@ import rowers.stravastuff as stravastuff from rowers.dataprep import rdata import rowers.dataprep as dataprep import rowers.metrics as metrics +import rowers.c2stuff as c2stuff from rowers.metrics import axes,axlabels,yaxminima,yaxmaxima @@ -1815,7 +1816,7 @@ def interactive_agegroupcpchart(age,normalized=False): fhpower = [] for distance in distances: - worldclasspower = metrics.getagegrouprecord( + worldclasspower = c2stuff.getagegrouprecord( age, sex='female', distance=distance, @@ -1829,7 +1830,7 @@ def interactive_agegroupcpchart(age,normalized=False): except ZeroDivisionError: pass for duration in durations: - worldclasspower = metrics.getagegrouprecord( + worldclasspower = c2stuff.getagegrouprecord( age, sex='female', duration=duration, @@ -1847,7 +1848,7 @@ def interactive_agegroupcpchart(age,normalized=False): flpower = [] for distance in distances: - worldclasspower = metrics.getagegrouprecord( + worldclasspower = c2stuff.getagegrouprecord( age, sex='female', distance=distance, @@ -1861,7 +1862,7 @@ def interactive_agegroupcpchart(age,normalized=False): except ZeroDivisionError: pass for duration in durations: - worldclasspower = metrics.getagegrouprecord( + worldclasspower = c2stuff.getagegrouprecord( age, sex='female', duration=duration, @@ -1879,7 +1880,7 @@ def interactive_agegroupcpchart(age,normalized=False): mlpower = [] for distance in distances: - worldclasspower = metrics.getagegrouprecord( + worldclasspower = c2stuff.getagegrouprecord( age, sex='male', distance=distance, @@ -1893,7 +1894,7 @@ def interactive_agegroupcpchart(age,normalized=False): except ZeroDivisionError: pass for duration in durations: - worldclasspower = metrics.getagegrouprecord( + worldclasspower = c2stuff.getagegrouprecord( age, sex='male', duration=duration, @@ -1912,7 +1913,7 @@ def interactive_agegroupcpchart(age,normalized=False): mhpower = [] for distance in distances: - worldclasspower = metrics.getagegrouprecord( + worldclasspower = c2stuff.getagegrouprecord( age, sex='male', distance=distance, @@ -1926,7 +1927,7 @@ def interactive_agegroupcpchart(age,normalized=False): except ZeroDivisionError: pass for duration in durations: - worldclasspower = metrics.getagegrouprecord( + worldclasspower = c2stuff.getagegrouprecord( age, sex='male', duration=duration, @@ -4060,15 +4061,7 @@ def thumbnails_set(r,id,favorites): columns += [f.yparam2 for f in favorites] columns += ['time'] - try: - rowdata = dataprep.getsmallrowdata_db(columns,ids=[id],doclean=True) - except: - return [ - {'script':"", - 'div':"", - 'notes':"" - }] - + rowdata = dataprep.getsmallrowdata_db(columns,ids=[id],doclean=True) rowdata.dropna(axis=1,how='all',inplace=True) diff --git a/rowers/metrics.py b/rowers/metrics.py index dd1852b2..1a488afa 100644 --- a/rowers/metrics.py +++ b/rowers/metrics.py @@ -6,7 +6,7 @@ from __future__ import unicode_literals from __future__ import absolute_import from rowers.utils import lbstoN import numpy as np -from rowers.models import C2WorldClassAgePerformance + import pandas as pd from scipy import optimize from django.utils import timezone @@ -290,8 +290,14 @@ rowingmetrics = ( ) - - +dtypes = {} + +for name,d in rowingmetrics: + if d['numtype'] == 'float': + dtypes[name] = float + elif d['numtype'] == 'int': + dtypes[name] = int + axesnew = [ (name,d['verbose_name'],d['ax_min'],d['ax_max'],d['type']) for name,d in rowingmetrics ] @@ -391,64 +397,3 @@ def calc_trimp(df,sex,hrmax,hrmin,hrftp): return trimp,hrtss -def getagegrouprecord(age,sex='male',weightcategory='hwt', - distance=2000,duration=None,indf=pd.DataFrame()): - - if not indf.empty: - if not duration: - df = indf[indf['distance'] == distance] - else: - duration = 60*int(duration) - df = indf[indf['duration'] == duration] - else: - if not duration: - df = pd.DataFrame( - list( - C2WorldClassAgePerformance.objects.filter( - distance=distance, - sex=sex, - weightcategory=weightcategory - ).values() - ) - ) - else: - duration=60*int(duration) - df = pd.DataFrame( - list( - C2WorldClassAgePerformance.objects.filter( - duration=duration, - sex=sex, - weightcategory=weightcategory - ).values() - ) - ) - - if not df.empty: - ages = df['age'] - powers = df['power'] - - #poly_coefficients = np.polyfit(ages,powers,6) - fitfunc = lambda pars, x: np.abs(pars[0])*(1-x/max(120,pars[1]))-np.abs(pars[2])*np.exp(-x/np.abs(pars[3]))+np.abs(pars[4])*(np.sin(np.pi*x/max(50,pars[5]))) - errfunc = lambda pars, x,y: fitfunc(pars,x)-y - - p0 = [700,120,700,10,100,100] - - try: - p1, success = optimize.leastsq(errfunc,p0[:], - args = (ages,powers)) - except: - p1 = p0 - success = 0 - - if success: - power = fitfunc(p1, float(age)) - - #power = np.polyval(poly_coefficients,age) - - power = 0.5*(np.abs(power)+power) - else: - power = 0 - else: - power = 0 - - return power diff --git a/rowers/models.py b/rowers/models.py index 02fb531f..defc99f8 100644 --- a/rowers/models.py +++ b/rowers/models.py @@ -28,6 +28,8 @@ from django_countries.fields import CountryField from scipy.interpolate import splprep, splev, CubicSpline import numpy as np +import shutil + from django.conf import settings from sqlalchemy import create_engine import sqlalchemy as sa @@ -2805,6 +2807,12 @@ def auto_delete_file_on_delete(sender, instance, **kwargs): if instance.csvfilename+'.gz': if os.path.isfile(instance.csvfilename+'.gz'): os.remove(instance.csvfilename+'.gz') + # remove parquet file + try: + dirname = 'media/strokedata_{id}.parquet.gz'.format(id=instance.id) + shutil.rmtree(dirname) + except FileNotFoundError: + pass @receiver(models.signals.post_delete,sender=Workout) def update_duplicates_on_delete(sender, instance, **kwargs): @@ -2842,20 +2850,20 @@ def update_duplicates_on_delete(sender, instance, **kwargs): # Delete stroke data from the database when a workout is deleted -@receiver(models.signals.post_delete,sender=Workout) -def auto_delete_strokedata_on_delete(sender, instance, **kwargs): - if instance.id: - query = sa.text('DELETE FROM strokedata WHERE workoutid={id};'.format( - id=instance.id, - )) - engine = create_engine(database_url, echo=False) - with engine.connect() as conn, conn.begin(): - try: - result = conn.execute(query) - except: - print("Database Locked") - conn.close() - engine.dispose() +#@receiver(models.signals.post_delete,sender=Workout) +#def auto_delete_strokedata_on_delete(sender, instance, **kwargs): +# if instance.id: +# query = sa.text('DELETE FROM strokedata WHERE workoutid={id};'.format( +# id=instance.id, +# )) +# engine = create_engine(database_url, echo=False) +# with engine.connect() as conn, conn.begin(): +# try: +# result = conn.execute(query) +# except: +# print("Database Locked") +# conn.close() +# engine.dispose() # Virtual Race results (for keeping results when workouts are deleted) @python_2_unicode_compatible diff --git a/rowers/serializers.py b/rowers/serializers.py index 7a99f2a1..d5ea3e84 100644 --- a/rowers/serializers.py +++ b/rowers/serializers.py @@ -7,7 +7,7 @@ from __future__ import unicode_literals # Also optionally define POST, PATCH methods (create, update) from rest_framework import serializers -from rowers.models import Workout,Rower,StrokeData,FavoriteChart +from rowers.models import Workout,Rower,FavoriteChart import datetime diff --git a/rowers/urls.py b/rowers/urls.py index b6bcc91b..cabc1753 100644 --- a/rowers/urls.py +++ b/rowers/urls.py @@ -7,7 +7,7 @@ from django.conf.urls import url, include from django.urls import path, re_path from django.contrib.auth.models import User from django.contrib.auth.decorators import login_required, permission_required -from rowers.models import Workout,Rower,StrokeData,FavoriteChart +from rowers.models import Workout,Rower,FavoriteChart from rest_framework import routers, serializers, viewsets,permissions from rest_framework.urlpatterns import format_suffix_patterns diff --git a/rowers/views/statements.py b/rowers/views/statements.py index 85df6cb6..069bd438 100644 --- a/rowers/views/statements.py +++ b/rowers/views/statements.py @@ -99,7 +99,7 @@ from rowers.models import ( ) from rowers.models import ( RowerPowerForm,RowerForm,GraphImage,AdvancedWorkoutForm, - RowerPowerZonesForm,AccountRowerForm,UserForm,StrokeData, + RowerPowerZonesForm,AccountRowerForm,UserForm, Team,TeamForm,TeamInviteForm,TeamInvite,TeamRequest, WorkoutComment,WorkoutCommentForm,RowerExportForm, CalcAgePerformance, @@ -256,19 +256,17 @@ def getfavorites(r,row): if 'speedcoach2' in row.workoutsource: workoutsource = 'speedcoach2' - try: - favorites = FavoriteChart.objects.filter(user=r, - workouttype__in=matchworkouttypes).order_by("id") - favorites2 = FavoriteChart.objects.filter(user=r, - workouttype__in=[workoutsource]).order_by("id") + favorites = FavoriteChart.objects.filter(user=r, + workouttype__in=matchworkouttypes).order_by("id") + favorites2 = FavoriteChart.objects.filter(user=r, + workouttype__in=[workoutsource]).order_by("id") - favorites = favorites | favorites2 - - maxfav = len(favorites)-1 - except: - favorites = None - maxfav = 0 + + favorites = favorites | favorites2 + + maxfav = len(favorites)-1 + return favorites,maxfav diff --git a/rowers/views/workoutviews.py b/rowers/views/workoutviews.py index 0293bfea..f8c0efbf 100644 --- a/rowers/views/workoutviews.py +++ b/rowers/views/workoutviews.py @@ -2767,7 +2767,6 @@ def workout_workflow_view(request,id): aantalcomments = len(comments) favorites,maxfav = getfavorites(r,row) - charts = get_call()