added extra data flushing event with separate trigger time
This commit is contained in:
+77
-42
@@ -14,45 +14,28 @@ from database import connect
|
|||||||
config = ConfigParser()
|
config = ConfigParser()
|
||||||
config.read(Path(__file__).parent / 'defaults.ini')
|
config.read(Path(__file__).parent / 'defaults.ini')
|
||||||
REPLICATES = int(config.get('main', 'replicates'))
|
REPLICATES = int(config.get('main', 'replicates'))
|
||||||
USER_ID = int(config.get('main', 'user_id'))
|
|
||||||
TRIGGER_INTERVAL_SECONDS = float(
|
TRIGGER_INTERVAL_SECONDS = float(
|
||||||
config.get(
|
config.get(
|
||||||
'main',
|
'main',
|
||||||
'trigger_interval_seconds'
|
'trigger_interval_seconds'
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
|
DATA_QUEUE_FLUSH_INTERVAL_SECONDS = float(
|
||||||
|
config.get(
|
||||||
|
'main',
|
||||||
|
'data_queue_flush_interval_seconds'
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def event(
|
def flush_data(
|
||||||
replicates: int = REPLICATES,
|
data_queue: Queue,
|
||||||
user_id: int = USER_ID
|
user_id: int
|
||||||
):
|
) -> None:
|
||||||
logging.debug('started event')
|
# connect to database
|
||||||
# load environment variables
|
|
||||||
replicates = int(os.getenv('REPLICATES', default=REPLICATES))
|
|
||||||
user_id = int(os.getenv('USER_ID', default=USER_ID))
|
|
||||||
# setup data queue
|
|
||||||
data_queue: Queue = Queue()
|
|
||||||
# run event
|
|
||||||
for rep_num in range(replicates):
|
|
||||||
logging.debug(f'running replicate {rep_num+1} of {replicates}')
|
|
||||||
# do measurement
|
|
||||||
try:
|
|
||||||
data = measure_bandwidth()
|
|
||||||
except Exception as e:
|
|
||||||
logging.error(f'failed to measure bandwidth: {e}')
|
|
||||||
continue
|
|
||||||
# add data to queue
|
|
||||||
data_queue.put(data)
|
|
||||||
# polite pause
|
|
||||||
time.sleep(1)
|
|
||||||
# check if queue should be flushed
|
|
||||||
queue_size = data_queue.qsize()
|
|
||||||
if queue_size == 0:
|
|
||||||
logging.debug('queue empty, continuing')
|
|
||||||
return
|
|
||||||
# flush queue to database
|
|
||||||
db = connect()
|
db = connect()
|
||||||
|
# flush queue
|
||||||
|
queue_size = data_queue.qsize()
|
||||||
for i in range(queue_size):
|
for i in range(queue_size):
|
||||||
# get data
|
# get data
|
||||||
data = data_queue.get()
|
data = data_queue.get()
|
||||||
@@ -78,26 +61,78 @@ def event(
|
|||||||
logging.debug('finished')
|
logging.debug('finished')
|
||||||
|
|
||||||
|
|
||||||
|
def event(
|
||||||
|
data_queue: Queue,
|
||||||
|
user_id: int,
|
||||||
|
replicates: int = REPLICATES
|
||||||
|
) -> None:
|
||||||
|
# run event
|
||||||
|
for rep_num in range(replicates):
|
||||||
|
logging.debug(f'running replicate {rep_num+1} of {replicates}')
|
||||||
|
# do measurement
|
||||||
|
try:
|
||||||
|
data = measure_bandwidth()
|
||||||
|
except Exception as e:
|
||||||
|
logging.error(f'failed to measure bandwidth: {e}')
|
||||||
|
continue
|
||||||
|
# add data to queue
|
||||||
|
data_queue.put(data)
|
||||||
|
# polite pause
|
||||||
|
time.sleep(1)
|
||||||
|
# flush queue
|
||||||
|
if data_queue.qsize() > 0:
|
||||||
|
flush_data(
|
||||||
|
data_queue=data_queue,
|
||||||
|
user_id=user_id
|
||||||
|
)
|
||||||
|
logging.debug('finished')
|
||||||
|
|
||||||
|
|
||||||
if __name__ == '__main__':
|
if __name__ == '__main__':
|
||||||
|
# setup
|
||||||
initialise_app()
|
initialise_app()
|
||||||
|
user_id = int(os.getenv('USER_ID', default=0))
|
||||||
|
replicates = int(os.getenv('REPLICATES', default=REPLICATES))
|
||||||
|
trigger_interval = float(
|
||||||
|
os.getenv(
|
||||||
|
'TRIGGER_INTERVAL_SECONDS',
|
||||||
|
default=TRIGGER_INTERVAL_SECONDS
|
||||||
|
)
|
||||||
|
)
|
||||||
|
data_flush_interval = float(
|
||||||
|
os.getenv(
|
||||||
|
'DATA_QUEUE_FLUSH_INTERVAL_SECONDS',
|
||||||
|
default=DATA_QUEUE_FLUSH_INTERVAL_SECONDS
|
||||||
|
)
|
||||||
|
)
|
||||||
|
data_queue: Queue = Queue()
|
||||||
|
# start main loop
|
||||||
trigger_time = time.time()
|
trigger_time = time.time()
|
||||||
|
data_flush_time = trigger_time + data_flush_interval
|
||||||
while True:
|
while True:
|
||||||
|
# measure event
|
||||||
if time.time() >= trigger_time:
|
if time.time() >= trigger_time:
|
||||||
logging.debug('triggered event')
|
logging.debug('triggered event')
|
||||||
# set new trigger time
|
# set new trigger time
|
||||||
trigger_interval = float(
|
|
||||||
os.getenv(
|
|
||||||
'TRIGGER_INTERVAL_SECONDS',
|
|
||||||
default=TRIGGER_INTERVAL_SECONDS
|
|
||||||
)
|
|
||||||
)
|
|
||||||
trigger_time += trigger_interval
|
trigger_time += trigger_interval
|
||||||
# run event
|
# run event
|
||||||
try:
|
event(
|
||||||
event()
|
data_queue=data_queue,
|
||||||
except Exception as e:
|
user_id=user_id,
|
||||||
logging.error(e)
|
replicates=replicates
|
||||||
else:
|
)
|
||||||
logging.info('finished event')
|
# extra data flush
|
||||||
|
if (
|
||||||
|
time.time() >= data_flush_time and
|
||||||
|
data_queue.qsize() > 0
|
||||||
|
):
|
||||||
|
logging.debug('triggered extra data flush')
|
||||||
|
# set new trigger time
|
||||||
|
data_flush_time += data_flush_interval
|
||||||
|
# flush data
|
||||||
|
flush_data(
|
||||||
|
data_queue=data_queue,
|
||||||
|
user_id=user_id
|
||||||
|
)
|
||||||
# polite pause
|
# polite pause
|
||||||
time.sleep(1)
|
time.sleep(1)
|
||||||
|
|||||||
Reference in New Issue
Block a user