使用扭曲中的反应器不断扫描数据库表时出现以下错误

发布于 2024-09-12 09:43:56 字数 1787 浏览 7 评论 0原文

成功运行几个小时后,我收到以下错误。

Traceback (most recent call last): 
  File "/usr/lib/python2.6/threading.py", line 484, in run 
    self.__target(*self.__args, **self.__kwargs) 
  File "/usr/lib/python2.6/dist-packages/twisted/python/threadpool.py", line 210, in _worker 
    result = context.call(ctx, function, *args, **kwargs) 
  File "/usr/lib/python2.6/dist-packages/twisted/python/context.py", line 59, in callWithContext 
    return self.currentContext().callWithContext(ctx, func, *args, **kw) 
  File "/usr/lib/python2.6/dist-packages/twisted/python/context.py", line 37, in callWithContext 
    return func(*args,**kw) 
--- <exception caught here> --- 
  File "/usr/lib/python2.6/dist-packages/twisted/enterprise/adbapi.py", line 436, in _runInteraction 
    conn.rollback() 
  File "/usr/lib/python2.6/dist-packages/twisted/enterprise/adbapi.py", line 52, in rollback 
    self._connection.rollback() 
_mysql_exceptions.OperationalError: (2006, 'MySQL server has gone away') 

我的代码是这样的...

from twisted.internet import reactor, defer,threads
from twisted.enterprise import adbapi 

dbpool = adbapi.ConnectionPool("MySQLdb", '192.168.1.102','test', 'test', 'test')
class Scanner:
 def _execQuery(self,txn):
      sql="SELECT tool_id,tool_name FROM tool_master"
      txn.execute(sql)
      result = txn.fetchall()
      return result
 def objCursor(self):
     return dbpool.runInteraction(self._execQuery)


 def printResult(self,result):
     print "resssssssssssssssssss",result
     reactor.callLater(3,self.deferExecute) 

 def deferExecute(self):
     self.objCursor().addCallback(self.printResult)
Scanner()   

class MyApp(object):
 reactor.callInThread(Scanner().deferExecute)
 reactor.run()
MyApp()

谁能告诉我为什么会收到此错误?

I am getting the following error after a few hours of successful running.

Traceback (most recent call last): 
  File "/usr/lib/python2.6/threading.py", line 484, in run 
    self.__target(*self.__args, **self.__kwargs) 
  File "/usr/lib/python2.6/dist-packages/twisted/python/threadpool.py", line 210, in _worker 
    result = context.call(ctx, function, *args, **kwargs) 
  File "/usr/lib/python2.6/dist-packages/twisted/python/context.py", line 59, in callWithContext 
    return self.currentContext().callWithContext(ctx, func, *args, **kw) 
  File "/usr/lib/python2.6/dist-packages/twisted/python/context.py", line 37, in callWithContext 
    return func(*args,**kw) 
--- <exception caught here> --- 
  File "/usr/lib/python2.6/dist-packages/twisted/enterprise/adbapi.py", line 436, in _runInteraction 
    conn.rollback() 
  File "/usr/lib/python2.6/dist-packages/twisted/enterprise/adbapi.py", line 52, in rollback 
    self._connection.rollback() 
_mysql_exceptions.OperationalError: (2006, 'MySQL server has gone away') 

My code is something like this...

from twisted.internet import reactor, defer,threads
from twisted.enterprise import adbapi 

dbpool = adbapi.ConnectionPool("MySQLdb", '192.168.1.102','test', 'test', 'test')
class Scanner:
 def _execQuery(self,txn):
      sql="SELECT tool_id,tool_name FROM tool_master"
      txn.execute(sql)
      result = txn.fetchall()
      return result
 def objCursor(self):
     return dbpool.runInteraction(self._execQuery)


 def printResult(self,result):
     print "resssssssssssssssssss",result
     reactor.callLater(3,self.deferExecute) 

 def deferExecute(self):
     self.objCursor().addCallback(self.printResult)
Scanner()   

class MyApp(object):
 reactor.callInThread(Scanner().deferExecute)
 reactor.run()
MyApp()

Can anyone tell me why I am getting this error?

如果你对这篇内容有疑问,欢迎到本站社区发帖提问 参与讨论,获取更多帮助,或者扫码二维码加入 Web 技术交流群。

扫码二维码加入Web技术交流群

发布评论

需要 登录 才能够评论, 你可以免费 注册 一个本站的账号。

评论(3

咆哮 2024-09-19 09:43:56

谁能告诉我为什么会出现这个错误..因为你做错了。

runInteraction 使用游标参数运行提供的函数到在线程中运行的事务。您不应该调用reactor.callInThread(Scanner().deferExecute)。

  • 最好使用twisted.internet.task.LoopingCall,它将确保调用在下一个调用之前完成。
  • 您只是在示例中运行查询,因此您可以仅使用 ConnectionPool.runQuery 而不是 ConnectionPool.runInteraction。
  • 使用 errorBack 函数报告异常。

尝试纠正格式错误的代码,我认为您已经得到了这个:

from twisted.internet import reactor, defer,threads
from twisted.enterprise import adbapi

dbpool = adbapi.ConnectionPool("MySQLdb", '192.168.1.102','test', 'test', 'test')

class Scanner:
  def _execQuery(self,txn):
    sql="SELECT tool_id,tool_name FROM tool_master"
    txn.execute(sql)
    result = txn.fetchall()
    return result

  def objCursor(self):
    return dbpool.runInteraction(self._execQuery)

  def printResult(self,result):
    print "resssssssssssssssssss",result
    reactor.callLater(3,self.deferExecute)

  def deferExecute(self):
    self.objCursor().addCallback(self.printResult)

Scanner()

class MyApp(object):
  reactor.callInThread(Scanner().deferExecute)
  reactor.run()

MyApp()

当您可能需要类似以下内容的时候。如果您计划编写扭曲的Application,则可以轻松修改此Scanner类以继承twisted.application.service.Service

from twisted.internet import reactor, defer, task
from twisted.enterprise import adbapi

class Scanner(object):
  def __init__(self,dbpool=None):
    self.dbpool = dbpool
    self.loopCall = task.LoopingCall(self.myQuery)
  def start(self):
    print "Started scanner"
    self.loopCall.start(3)
  def stop(self):
    print "Stopping scanner"
    self.loopCall.stop()
  def myQuery(self):
    def interact(txn):
      sql="SELECT tool_id,tool_name FROM tool_master"
      txn.execute(sql)
      return txn.fetchall()
    d = self.dbpool.runInteraction(interact)
    d.addCallbacks(self.printResult,self.printError)
  def printResult(self,result):
    print "Got Result: %r" % result
  def printError(self,error):
    print "Got Error: %r" % error
    error.printTraceback()

if __name__ == '__main__':
  from twisted.internet import reactor
  dbpool = adbapi.ConnectionPool("MySQLdb", '192.168.1.102','test', 'test', 'test')
  s = Scanner(dbpool)
  reactor.callWhenRunning(s.start)
  reactor.addSystemEventTrigger('before','shutdown',s.stop)
  reactor.run()

can anyone tell me why I am getting this error.. because you're doing it wrong.

runInteraction runs the supplied function with an argument of a cursor to a transaction which is run in a thread. You shouldn't be calling reactor.callInThread(Scanner().deferExecute).

  • It's better to use a twisted.internet.task.LoopingCall, it will make sure that the call completes before the next is fired.
  • You're just running a query in your example, so you could just use ConnectionPool.runQuery instead of ConnectionPool.runInteraction.
  • Use errorBack functions to report on Exceptions.

Attempting to correct for your badly formatted code, I think you've got this:

from twisted.internet import reactor, defer,threads
from twisted.enterprise import adbapi

dbpool = adbapi.ConnectionPool("MySQLdb", '192.168.1.102','test', 'test', 'test')

class Scanner:
  def _execQuery(self,txn):
    sql="SELECT tool_id,tool_name FROM tool_master"
    txn.execute(sql)
    result = txn.fetchall()
    return result

  def objCursor(self):
    return dbpool.runInteraction(self._execQuery)

  def printResult(self,result):
    print "resssssssssssssssssss",result
    reactor.callLater(3,self.deferExecute)

  def deferExecute(self):
    self.objCursor().addCallback(self.printResult)

Scanner()

class MyApp(object):
  reactor.callInThread(Scanner().deferExecute)
  reactor.run()

MyApp()

When you probably need something like the following instead. If you're planning on writing a twisted Application will be easy to modify this Scanner class to inherit from twisted.application.service.Service.

from twisted.internet import reactor, defer, task
from twisted.enterprise import adbapi

class Scanner(object):
  def __init__(self,dbpool=None):
    self.dbpool = dbpool
    self.loopCall = task.LoopingCall(self.myQuery)
  def start(self):
    print "Started scanner"
    self.loopCall.start(3)
  def stop(self):
    print "Stopping scanner"
    self.loopCall.stop()
  def myQuery(self):
    def interact(txn):
      sql="SELECT tool_id,tool_name FROM tool_master"
      txn.execute(sql)
      return txn.fetchall()
    d = self.dbpool.runInteraction(interact)
    d.addCallbacks(self.printResult,self.printError)
  def printResult(self,result):
    print "Got Result: %r" % result
  def printError(self,error):
    print "Got Error: %r" % error
    error.printTraceback()

if __name__ == '__main__':
  from twisted.internet import reactor
  dbpool = adbapi.ConnectionPool("MySQLdb", '192.168.1.102','test', 'test', 'test')
  s = Scanner(dbpool)
  reactor.callWhenRunning(s.start)
  reactor.addSystemEventTrigger('before','shutdown',s.stop)
  reactor.run()
平生欢 2024-09-19 09:43:56

经过所有建议& Matt 的帮助我有以下代码正在成功运行:

#!usr/bin/python
# Using the "dbmodule" from the previous example, create a ConnectionPool
from twisted.internet import reactor
from twisted.enterprise import adbapi
from twisted.internet import reactor, defer,threads
from twisted.python.threadpool import ThreadPool
import itertools
from twisted.internet.threads import deferToThread
from twisted.internet import reactor, defer, task
from tools.printTime import *
from tools.getVersion import *
from sh_log import *


concurrent = 30
finished=itertools.count(1)
reactor.suggestThreadPoolSize(concurrent)


#Creating Global Instance variables
path="tools"
lo=Log()

class ToolsBuilder:

  def build(self,txn,tool,asset_id):
    if tool:
        print "\n"
        try:
            sql="select tool_filename from tool_master where tool_id = %s" %(tool,)
            sql_asset="select asset_url from asset_master where asset_id = %s" %(asset_id,)
            txn.execute(sql_asset)
            asset_url = txn.fetchall()
            log_date=lo.log_date()
            txn.execute(sql)
            result = txn.fetchall()

            log='\n'+log_date+"::"+str(result[0][0])+ " tool object is created......\n"
            lo.wfile(log)
            temp=(path +'/' + str(result[0][0]))

            if result:
                if temp:
                    f=open(temp).read()
                    obj_tool=compile(f, 'a_filename', 'exec')
                    return obj_tool
        except:
            lo.wfile("Error in creating executable tool object......")
tb=ToolsBuilder()

class ToolsVectorGenerator:

  def generate(self,txn,res_set={}):

    v1=[]
    for asset_id in res_set.iterkeys():
        try:
            obj_tools=[]
            if asset_id:
                print "asset_id..............................",asset_id
                log_date=lo.log_date()
                log=log_date+"::"+" \nVector generation for the asset number...:"+str(asset_id)
                lo.wfile(log)
                vector=[]
                tools_arr=[]
                obj_tools=[]
                for tool in res_set[asset_id]:
                    if tool:
                        print "tool..............",tool
                        temp_tool=tb.build(txn,tool,asset_id)
                        print "temp_tool..........",temp_tool

                        #fetch data of tool setting.....
                        sql_tool_setting="select * from tool_asset_settings where tool_id =%s" %(tool,)
                        txn.execute(sql_tool_setting)
                        result_tool_setting = txn.fetchall()
                        tool_id=result_tool_setting[0][1]
                        t_id=int(tool_id)
                        tool_id_arr=[]
                        tool_id_arr.append(t_id)
                        tool_id_arr.append(result_tool_setting)
                        tool_id_arr.append(temp_tool)
                        tools_arr.append(tool_id_arr)

                #fetch data from asset master
                sql_asset="select asset_name from asset_master where asset_id=%s" %(asset_id,)
                txn.execute(sql_asset)
                result_asset = txn.fetchall()
                vector.append(result_asset)
                vector.append(tools_arr)
        except:
            lo.wfile("\nError in getting asset,please check your database or network connection......")
        tvm.executeVector(vector)
tvg=ToolsVectorGenerator()

class Tool:

  def exectool(self,tool):
      exec tool
      return

  def getResult(self,tool):
      return deferToThread(self.exectool, tool)

to=Tool()

class StateMachine:

  def setPriority(self,txn,tup):
        temp=[]
        arr=[]
        for li in tup:
            sql2="select tool_dependency from tool_asset_settings where tool_id =%s" %(li[1],)
            txn.execute(sql2)
            result12 = txn.fetchall()
            arr=[]
            if result12[0][0]!=None:
                tup12=result12[0][0]
                arr=(li[0],tup12)
    #                print "arr.........",arr
                if arr in tup:
                    print "This element is already exist......."
                else:
                    temp.append(arr)
        temp.extend(tup)
        return tuple(temp)

st=StateMachine()

class ToolsVectorExecutionManager(object):

    def executeVector(self,toolsvector):
        print "toolsvector================>",toolsvector
        if toolsvector:
            for tools in toolsvector[1]:
                if tools[2] != None:
                    to.getResult(tools[2])

tvm=ToolsVectorExecutionManager()

class ToolsToExecuteAnalyzer:

  def __init__(self,dbpool=None):
    self.dbpool = dbpool
    self.loopCall = task.LoopingCall(self.myQuery)

  def start(self):
    print "Started scanner"
    self.loopCall.start(3)

  def stop(self):
    print "Stopping scanner"
    self.loopCall.stop()

  def myQuery(self):

    def interact(txn):
      sql="SELECT tool_asset_id,tool_execute_id FROM tool_to_execute where status='0'"
      txn.execute(sql)
      result=txn.fetchall()
      if result:
            tool_asset_id=tuple([int(e[0]) for e in result])
            tool_execute_id=tuple([int(e[1]) for e in result])
            if len(tool_asset_id)>1:
                sql1="SELECT asset_id,tool_id FROM tool_in_assets WHERE tool_asset_id IN %s"%(tool_asset_id,)
            else:
                sql1="SELECT asset_id,tool_id FROM tool_in_assets WHERE tool_asset_id = (%s)"%(tool_asset_id)
            txn.execute(sql1)
            tup = txn.fetchall()

#dependency check for the selected tool
            asset_tool=st.setPriority(txn,tup)

            log_date=lo.log_date()
            log=log_date+"::priority have been set for the tools......\n"
            lo.wfile(log)

#creating group of asset with their tools
            res={}
            for element in asset_tool:
                if element[0] in res:
                    res[element[0]].append(int(element[1]))
                else:
                    res[int(element[0])] = [int(element[1])]
#Recored deletion from tool_to_execute table
            if res!=None and res.keys()!=[]:
                for asset_id in res.iterkeys():
                    if len(tool_execute_id)>1:
                        sql_del="delete from tool_to_execute where tool_execute_id in %s " %(tool_execute_id,)
                    else:
                        sql_del="delete from tool_to_execute where tool_execute_id = %s" %(tool_execute_id)
                txn.execute(sql_del)

#New Addition of vector
            tvg.generate(txn,res)
#            return res

    d = self.dbpool.runInteraction(interact)
    d.addCallbacks(self.printResult,self.printError)

  def printResult(self,res):

    print "In printResult after generate...."

  def printError(self,error):

    print "Got Error: %r" % error
    error.printTraceback()

ToolsToExecuteAnalyzer()


if __name__ == '__main__':

  from twisted.internet import reactor

  dbpool = adbapi.ConnectionPool("MySQLdb", 'localhost', 'test', 'test','test')
  s = ToolsToExecuteAnalyzer(dbpool)
  reactor.callWhenRunning(s.start)
  reactor.addSystemEventTrigger('before','shutdown',s.stop)
  reactor.run()

这是我的整个代码,我只是想知道运行了多少个线程,对于每个工具来说意味着新线程?
无论如何,感谢马特的帮助..:)

After all the suggestion & help by Matt I have following code which is running successfully:

#!usr/bin/python
# Using the "dbmodule" from the previous example, create a ConnectionPool
from twisted.internet import reactor
from twisted.enterprise import adbapi
from twisted.internet import reactor, defer,threads
from twisted.python.threadpool import ThreadPool
import itertools
from twisted.internet.threads import deferToThread
from twisted.internet import reactor, defer, task
from tools.printTime import *
from tools.getVersion import *
from sh_log import *


concurrent = 30
finished=itertools.count(1)
reactor.suggestThreadPoolSize(concurrent)


#Creating Global Instance variables
path="tools"
lo=Log()

class ToolsBuilder:

  def build(self,txn,tool,asset_id):
    if tool:
        print "\n"
        try:
            sql="select tool_filename from tool_master where tool_id = %s" %(tool,)
            sql_asset="select asset_url from asset_master where asset_id = %s" %(asset_id,)
            txn.execute(sql_asset)
            asset_url = txn.fetchall()
            log_date=lo.log_date()
            txn.execute(sql)
            result = txn.fetchall()

            log='\n'+log_date+"::"+str(result[0][0])+ " tool object is created......\n"
            lo.wfile(log)
            temp=(path +'/' + str(result[0][0]))

            if result:
                if temp:
                    f=open(temp).read()
                    obj_tool=compile(f, 'a_filename', 'exec')
                    return obj_tool
        except:
            lo.wfile("Error in creating executable tool object......")
tb=ToolsBuilder()

class ToolsVectorGenerator:

  def generate(self,txn,res_set={}):

    v1=[]
    for asset_id in res_set.iterkeys():
        try:
            obj_tools=[]
            if asset_id:
                print "asset_id..............................",asset_id
                log_date=lo.log_date()
                log=log_date+"::"+" \nVector generation for the asset number...:"+str(asset_id)
                lo.wfile(log)
                vector=[]
                tools_arr=[]
                obj_tools=[]
                for tool in res_set[asset_id]:
                    if tool:
                        print "tool..............",tool
                        temp_tool=tb.build(txn,tool,asset_id)
                        print "temp_tool..........",temp_tool

                        #fetch data of tool setting.....
                        sql_tool_setting="select * from tool_asset_settings where tool_id =%s" %(tool,)
                        txn.execute(sql_tool_setting)
                        result_tool_setting = txn.fetchall()
                        tool_id=result_tool_setting[0][1]
                        t_id=int(tool_id)
                        tool_id_arr=[]
                        tool_id_arr.append(t_id)
                        tool_id_arr.append(result_tool_setting)
                        tool_id_arr.append(temp_tool)
                        tools_arr.append(tool_id_arr)

                #fetch data from asset master
                sql_asset="select asset_name from asset_master where asset_id=%s" %(asset_id,)
                txn.execute(sql_asset)
                result_asset = txn.fetchall()
                vector.append(result_asset)
                vector.append(tools_arr)
        except:
            lo.wfile("\nError in getting asset,please check your database or network connection......")
        tvm.executeVector(vector)
tvg=ToolsVectorGenerator()

class Tool:

  def exectool(self,tool):
      exec tool
      return

  def getResult(self,tool):
      return deferToThread(self.exectool, tool)

to=Tool()

class StateMachine:

  def setPriority(self,txn,tup):
        temp=[]
        arr=[]
        for li in tup:
            sql2="select tool_dependency from tool_asset_settings where tool_id =%s" %(li[1],)
            txn.execute(sql2)
            result12 = txn.fetchall()
            arr=[]
            if result12[0][0]!=None:
                tup12=result12[0][0]
                arr=(li[0],tup12)
    #                print "arr.........",arr
                if arr in tup:
                    print "This element is already exist......."
                else:
                    temp.append(arr)
        temp.extend(tup)
        return tuple(temp)

st=StateMachine()

class ToolsVectorExecutionManager(object):

    def executeVector(self,toolsvector):
        print "toolsvector================>",toolsvector
        if toolsvector:
            for tools in toolsvector[1]:
                if tools[2] != None:
                    to.getResult(tools[2])

tvm=ToolsVectorExecutionManager()

class ToolsToExecuteAnalyzer:

  def __init__(self,dbpool=None):
    self.dbpool = dbpool
    self.loopCall = task.LoopingCall(self.myQuery)

  def start(self):
    print "Started scanner"
    self.loopCall.start(3)

  def stop(self):
    print "Stopping scanner"
    self.loopCall.stop()

  def myQuery(self):

    def interact(txn):
      sql="SELECT tool_asset_id,tool_execute_id FROM tool_to_execute where status='0'"
      txn.execute(sql)
      result=txn.fetchall()
      if result:
            tool_asset_id=tuple([int(e[0]) for e in result])
            tool_execute_id=tuple([int(e[1]) for e in result])
            if len(tool_asset_id)>1:
                sql1="SELECT asset_id,tool_id FROM tool_in_assets WHERE tool_asset_id IN %s"%(tool_asset_id,)
            else:
                sql1="SELECT asset_id,tool_id FROM tool_in_assets WHERE tool_asset_id = (%s)"%(tool_asset_id)
            txn.execute(sql1)
            tup = txn.fetchall()

#dependency check for the selected tool
            asset_tool=st.setPriority(txn,tup)

            log_date=lo.log_date()
            log=log_date+"::priority have been set for the tools......\n"
            lo.wfile(log)

#creating group of asset with their tools
            res={}
            for element in asset_tool:
                if element[0] in res:
                    res[element[0]].append(int(element[1]))
                else:
                    res[int(element[0])] = [int(element[1])]
#Recored deletion from tool_to_execute table
            if res!=None and res.keys()!=[]:
                for asset_id in res.iterkeys():
                    if len(tool_execute_id)>1:
                        sql_del="delete from tool_to_execute where tool_execute_id in %s " %(tool_execute_id,)
                    else:
                        sql_del="delete from tool_to_execute where tool_execute_id = %s" %(tool_execute_id)
                txn.execute(sql_del)

#New Addition of vector
            tvg.generate(txn,res)
#            return res

    d = self.dbpool.runInteraction(interact)
    d.addCallbacks(self.printResult,self.printError)

  def printResult(self,res):

    print "In printResult after generate...."

  def printError(self,error):

    print "Got Error: %r" % error
    error.printTraceback()

ToolsToExecuteAnalyzer()


if __name__ == '__main__':

  from twisted.internet import reactor

  dbpool = adbapi.ConnectionPool("MySQLdb", 'localhost', 'test', 'test','test')
  s = ToolsToExecuteAnalyzer(dbpool)
  reactor.callWhenRunning(s.start)
  reactor.addSystemEventTrigger('before','shutdown',s.stop)
  reactor.run()

This is my whole code, I just wanted to know how many threads running, means for each tool new thread?
Anyway, thanks Matt for your help..:)

何以笙箫默 2024-09-19 09:43:56

您可能还想看看这个片段,它提供了一个 ConnectionPool 子类,可以在“MySQL 服务器已消失”时重新连接。

http://www.gelens.org/2009/09/13/扭曲连接池重访/

You may also want to take a look at this snippet which provides a ConnectionPool subclass that reconnects on "MySQL server has gone away".

http://www.gelens.org/2009/09/13/twisted-connectionpool-revisited/

~没有更多了~
我们使用 Cookies 和其他技术来定制您的体验包括您的登录状态等。通过阅读我们的 隐私政策 了解更多相关信息。 单击 接受 或继续使用网站,即表示您同意使用 Cookies 和您的相关数据。
原文