You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用JPype连接读取AMQ队列时偶发挂起问题排查求助

问题:JPype连接AMQ Broker频繁挂起排查

我们使用JPype连接AMQ Broker并读取队列,程序运行基本正常,但频繁在连接状态下出现挂起。已排查确认防火墙无异常,现将使用的代码及版本信息提供如下,恳请协助排查问题。

代码实现

import jpype
import time
from datetime import datetime, timedelta

def consume_message(queue_name):
    try:
        print("Start the JVM")
        if not jpype.isJVMStarted():
            jpype.startJVM(jpype.getDefaultJVMPath())

        print("Set SSL properties")
        ssl_options = {
            "java.naming.security.protocol": "ssl",
            "javax.net.ssl.keyStore": args.keystore,
            "javax.net.ssl.keyStorePassword": args.keystore_pwd,
            "javax.net.ssl.trustStore": args.truststore,
            "javax.net.ssl.trustStorePassword": args.truststore_pwd
        }

        print("Create a connection factory")
        connectionFactory = jpype.JClass("org.apache.activemq.ActiveMQConnectionFactory")()
        connectionFactory.setBrokerURL(args.broker_url)

        print("Set SSL options")
        for key, value in ssl_options.items():
            jpype.java.lang.System.setProperty(key, value)

        print("Create a connection")
        connection = connectionFactory.createConnection()
        print("client_id")
        print(args.client_id)
        # connection.setClientID(args.client_id)
        connection.setDefaultClientID('queue' + str(datetime.now()))

        print("start")
        connection.start()

        print("Create session")
        session = connection.createSession(False, jpype.JInt(1))

        print("Lookup queue")
        queue = session.createQueue(queue_name)

        print("Create consumer")
        consumer = session.createConsumer(queue)

        print("Hello connecting to AMQ and started receiving messages.......................")
        print("Start receiving messages and wait until end time")
        end_time = datetime.now() + timedelta(minutes=args.end_time)

        while datetime.now() <= end_time:
            message = consumer.receiveNoWait()
            if message:
                # print(message)
                parse_message(message.getText())
                time.sleep(args.sleep_time)

        print("Done waiting......")

    except jpype.JException as e:
        print(f"JPype Exception: {e}")
    except Exception as e:
        print(f"An Error Occurred: {e}")
        return "An error occurred: " + str(e)
    finally:
        # Close connections and shutdown JVM
        if connection:
            connection.close()
        if jpype.isJVMStarted():
            jpype.shutdownJVM()

if __name__ == "__main__":
    if not jpype.isJVMStarted():
        jpype.startJVM(
            jpype.getDefaultJVMPath(),
            "-Djava.class.path=" + "/mtproject/activemq-all-5.10.2.jar",
            "-Djava.util.logging.config.file=" + "/mtproject/logging.properties"
        )
    
    print("start amq client...")
    queue_name = args.queue_name
    consume_message(queue_name)
    print("end amq client...")
    jpype.shutdownJVM()

使用版本

  • Jpype 1.3.0
  • AMQ 7.10.3.CR2-redhat-00001
  • Python 3.6.8

内容的提问来源于stack exchange,提问作者Satya

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.18 09:25:53