Skip to content

Commit ad2290d

Browse files
committed
Enhance error logging
1 parent 7ae296d commit ad2290d

File tree

1 file changed

+5
-3
lines changed

1 file changed

+5
-3
lines changed

hkube_python_wrapper/wrapper/algorunner.py

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -315,6 +315,7 @@ def _getMethod(self, name):
315315
if (self._algorithm):
316316
return self._algorithm.get(name)
317317
return None
318+
318319
def _init(self, options):
319320
redirector = None
320321
try:
@@ -345,14 +346,14 @@ def _init(self, options):
345346

346347
def _discovery_update(self, discovery):
347348
log.debug('Got discovery update {discovery}', discovery=discovery)
348-
messageListenerConfig = {'encoding': config.discovery['encoding'],'delay':config.discovery['delay']}
349+
messageListenerConfig = {'encoding': config.discovery['encoding'], 'delay': config.discovery['delay']}
349350
self.streamingManager.setupStreamingListeners(
350351
messageListenerConfig, discovery, self._job.nodeName)
351352

352353
def _setupStreamingProducer(self, nodeName):
353354
def onStatistics(statistics):
354355
thread_list = ""
355-
self._printThread = self._printThread + 1
356+
self._printThread = self._printThread + 1
356357
for thread in threading.enumerate():
357358
thread_list = thread_list + " " + str(thread.name)
358359
if (self._printThread % 30 == 0):
@@ -492,6 +493,7 @@ def _stopAlgorithm(self, options):
492493
if (self._job and self._job.isStreaming):
493494
if (forceStop is False):
494495
stoppingState = True
496+
495497
def stopping():
496498
while (stoppingState):
497499
self._sendCommand(messages.outgoing.stopping, None)
@@ -559,7 +561,7 @@ def _sendCommand(self, command, data):
559561

560562
def sendError(self, error):
561563
try:
562-
log.error(error)
564+
log.error("Sending error to worker " + str(error))
563565
self._wsc.send({
564566
'command': messages.outgoing.error,
565567
'error': {

0 commit comments

Comments
 (0)