Changeset: 0e1d8ffb5727 for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=0e1d8ffb5727
Modified Files:
        sql/test/remote/Tests/creds.SQL.py.in
Branch: remote_auth
Log Message:

Heavily refactor the creds test to make it easier to reuse


diffs (125 lines):

diff --git a/sql/test/remote/Tests/creds.SQL.py.in 
b/sql/test/remote/Tests/creds.SQL.py.in
--- a/sql/test/remote/Tests/creds.SQL.py.in
+++ b/sql/test/remote/Tests/creds.SQL.py.in
@@ -30,6 +30,17 @@ RATINGS_TABLE_DEF = ''' (
 )
 '''
 
+# Complicated FK constraints on merge/remote tables do not work great
+# currently (May 2018)
+RATINGS_TABLE_DEF_FK = ''' (
+    movie_id BIGINT,
+    customer_id BIGINT,
+    rating TINYINT,
+    rating_date DATE
+)
+'''
+
+# Find a free network port
 def freeport():
     sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
     sock.bind(('', 0))
@@ -37,54 +48,67 @@ def freeport():
     sock.close()
     return port
 
-def worker_load(workerrec):
-    
filename="$TSTDATAPATH/netflix_data/ratings_sample_{}.csv".format(workerrec['num'])
+# Create the remote tables and load the data. Note: the supervisor
+# database should be already started and should contain the movies
+# table.
+def worker_load(in_filename, workerrec, cmovies, ratings_table_def_fk):
     c = workerrec['conn']
-    cmovies = "CREATE REMOTE TABLE movies {} ON '{}'".format(MOVIES_TABLE_DEF, 
supervisor_uri)
-    screateq = "CREATE TABLE ratings {}".format(RATINGS_TABLE_DEF)
-    load_data = "COPY INTO ratings FROM '{}' USING DELIMITERS 
',','\n'".format(filename)
+    screateq = "CREATE TABLE ratings {}".format(ratings_table_def_fk)
+    load_data = "COPY INTO ratings FROM '{}' USING DELIMITERS 
',','\n'".format(in_filename)
     c.execute(cmovies)
     c.execute(screateq)
     c.execute(load_data)
 
+# Setup and start workers
+def create_workers(fn_template, nworkers, cmovies, ratings_table_def_fk):
+    workers = []
+    for i in range(nworkers):
+        workerport = freeport()
+        workerdbname = 'worker_{}'.format(i)
+        workerrec = {
+            'num': i,
+            'port': workerport,
+            'dbname': workerdbname,
+            'dbfarm': os.path.join(TMPDIR, workerdbname),
+            'mapi': 
'mapi:monetdb://localhost:{}/{}/sys/ratings'.format(workerport, workerdbname),
+        }
+        workerrec['proc'] = process.server(mapiport=workerrec['port'], 
dbname=workerrec['dbname'], dbfarm=workerrec['dbfarm'], stdin=process.PIPE, 
stdout=process.PIPE)
+        workerrec['conn'] = pymonetdb.connect(database=workerrec['dbname'], 
port=workerport, autocommit=True)
+        filename = fn_template.format(workerrec['num'])
+        t = threading.Thread(target=worker_load, args=[filename, workerrec, 
cmovies, ratings_table_def_fk])
+        t.start()
+        workerrec['loadthread'] = t
+        workers.append(workerrec)
+
+    for wrec in workers:
+        wrec['loadthread'].join()
+
+    return workers
+
+# Start supervisor database
 supervisorport = freeport()
 supervisorproc = process.server(mapiport=supervisorport, dbname="supervisor", 
dbfarm=os.path.join(TMPDIR, "supervisor"), stdin=process.PIPE, 
stdout=process.PIPE)
 supervisorconn = pymonetdb.connect(database='supervisor', port=supervisorport, 
autocommit=True)
 supervisor_uri = 
"mapi:monetdb://localhost:{}/supervisor".format(supervisorport)
 c = supervisorconn.cursor()
 
+# Create the movies table and load the data
 movies_filename="$TSTDATAPATH/netflix_data/movies.csv"
 movies_create = "CREATE TABLE movies {}".format(MOVIES_TABLE_DEF)
 c.execute(movies_create)
 load_movies = "COPY INTO movies FROM '{}' USING DELIMITERS 
',','\n','\"'".format(movies_filename)
 c.execute(load_movies)
 
-# Setup and start workers
-workers = []
-for i in range(NWORKERS):
-    workerport = freeport()
-    workerdbname = 'worker_{}'.format(i)
-    workerrec = {
-        'num': i,
-        'port': workerport,
-        'dbname': workerdbname,
-        'dbfarm': os.path.join(TMPDIR, workerdbname),
-        'mapi': 
'mapi:monetdb://localhost:{}/{}/sys/ratings'.format(workerport, workerdbname),
-    }
-    workerrec['proc'] = process.server(mapiport=workerrec['port'], 
dbname=workerrec['dbname'], dbfarm=workerrec['dbfarm'], stdin=process.PIPE, 
stdout=process.PIPE)
-    workerrec['conn'] = pymonetdb.connect(database=workerrec['dbname'], 
port=workerport, autocommit=True)
-    t = threading.Thread(target=worker_load, args=[workerrec])
-    t.start()
-    workerrec['loadthread'] = t
-    workers.append(workerrec)
-
-for wrec in workers:
-    wrec['loadthread'].join()
-
-
+# Declare the ratings merge table on supervisor
 mtable = "CREATE MERGE TABLE ratings {}".format(RATINGS_TABLE_DEF)
 c.execute(mtable)
 
+# Create the workers and load the ratings data
+fn_template="$TSTDATAPATH/netflix_data/ratings_sample_{}.csv"
+cmovies = "CREATE REMOTE TABLE movies {} ON '{}' WITH USER 'monetdb' PASSWORD 
'monetdb'".format(MOVIES_TABLE_DEF, supervisor_uri)
+workers = create_workers(fn_template, NWORKERS, cmovies, RATINGS_TABLE_DEF_FK)
+
+# Create the remote tables on supervisor
 for wrec in workers:
     rtable = "CREATE REMOTE TABLE ratings{} {} on '{}' WITH USER 'monetdb' 
PASSWORD 'monetdb'".format(wrec['num'], RATINGS_TABLE_DEF, wrec['mapi'])
     c.execute(rtable)
@@ -92,6 +116,7 @@ for wrec in workers:
     atable = "ALTER TABLE ratings add table ratings{}".format(wrec['num'])
     c.execute(atable)
 
+# Run the queries
 c.execute("SELECT COUNT(*) FROM ratings0")
 print("{} rows in remote table".format(c.fetchall()[0][0]))
 
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list

Reply via email to