【问题标题】:Julia ZMQ - connecting to other WebSockets produces StateErrorJulia ZMQ - 连接到其他 WebSocket 会产生 StateError
【发布时间】:2021-11-18 23:49:35
【问题描述】:

我正在尝试使用 ZMQ 将多个发布者连接到一个订阅者 (python)。这是一个这样的发布者(我使用连接而不是绑定,因为订阅者绑定)。在我取消阻止下面的注释代码之前,代码可以正常工作。

然后我在 Windows 上收到此错误:

LoadError: StateError("Unknown error")

Ubuntu 上:

StateError("Socket operation on non-socket")

源代码:

using ZMQ
using WebSockets
using JSON3

const uri = "wss://ws.okex.com:8443/ws/v5/public"

function produce_string()
    return "hi"
end

function main()
    payload = Dict(
            :op => "subscribe",
            :args => [
                Dict(
                    "channel" => "books50-l2-tbt",
                    "instType" => "Futures",
                    "instId" => "FIL-USD-220325",
                ),
            ],
        )
    # Unblock this code to produce error
    # @async while true
    #     WebSockets.open(uri) do ws
    #         confirmation = true
    #         if isopen(ws)
    #             write(ws, JSON3.write(payload))
    #         end
    #     end
    # end

    ctx = Context()
    zmq_socket = Socket(ctx, PUB)
    addr = "tcp://localhost:" * string(8093)
    ZMQ.connect(zmq_socket, addr)
    sleep(3)
    ZMQ.send(zmq_socket, "hi")

    while true
        my_string = produce_string()
        ZMQ.send(zmq_socket, my_string)
        println("sent")
        sleep(1)
    end

end

main()

【问题讨论】:

    标签: tcp julia ipc zeromq pyzmq


    【解决方案1】:

    这似乎至少部分是一个错误(或难以理解的行为),所以我建议你在 repo 上创建一个问题。可能与:Test Error: Assertion failed: Socket operation on non-socket #147有关。

    但是,我们可以尽最大努力尝试了解问题所在,或许可以找到解决方法。由于 ZMQ.jl 使用 libzmq 在低级别处理套接字,它可能会干扰 Julia 对文件描述符的处理,我们可能有一个race condition。让我们通过稍微修改一下代码来检验这个假设:

        @async WebSockets.open(uri) do ws
            while true
                if isopen(ws)
                    msg = JSON3.write(payload)
                    write(ws, msg)
                    display(ws.socket.bio)
                    break
                end
            end
        end
    
        sleep(0.1)
        ctx = Context()
        zmq_socket = Socket(ctx, PUB)
        dump(zmq_socket)
        addr = "tcp://localhost:" * string(8093)
        ZMQ.connect(zmq_socket, addr)
        sleep(3)
        ZMQ.send(zmq_socket, "hi")
    

    我只是更改了一些东西以获取打印出必要信息的代码。我们看到:

    Socket
      data: Ptr{Nothing} @0x0000000001d5e590
      pollfd: FileWatching._FDWatcher
        handle: Ptr{Nothing} @0x00000000018b7970
        fdnum: Int64 31
        refcount: Tuple{Int64, Int64}
          1: Int64 1
          2: Int64 0
        notify: Base.GenericCondition{Base.Threads.SpinLock}
          waitq: Base.InvasiveLinkedList{Task}
            head: Nothing nothing
            tail: Nothing nothing
          lock: Base.Threads.SpinLock
            owned: Int64 0
        events: Int32 0
        active: Tuple{Bool, Bool}
          1: Bool false
          2: Bool false
    

    TCPSocket(RawFD(31) paused, 0 bytes waiting)
    

    pollfd.fdnum 字段是 31,这与 TCPSocket 文件描述符相同,所以可能就是这样。

    我们能做什么?

    在上面的代码中,我已经对您的原始代码进行了更改,我将调用中的while 循环移动到了WebSockets.open,您真的要在每个循环中打开一个新套接字吗?其次,我们可以尝试同步一下我们的线程,以确保在调用 ZMQ 之前我们已经完成了套接字的打开:

    function main()
        payload = Dict(
                :op => "subscribe",
                :args => [
                    Dict(
                        "channel" => "books50-l2-tbt",
                        "instType" => "Futures",
                        "instId" => "FIL-USD-220325",
                    ),
                ],
            )
        msg_channel = Channel(1)
        @async WebSockets.open(uri) do ws
            while true
                if isopen(ws)
                    msg = JSON3.write(payload)
                    put!(msg_channel, msg)
                    write(ws, msg)
                end
            end
        end
    
        println(take!(msg_channel))
        ctx = Context()
        zmq_socket = Socket(ctx, PUB)
        addr = "tcp://localhost:" * string(8093)
        ZMQ.connect(zmq_socket, addr)
        sleep(3)
        ZMQ.send(zmq_socket, "hi")
    
        while true
            my_string = produce_string()
            ZMQ.send(zmq_socket, my_string)
            println("sent")
            sleep(1)
        end
    end
    

    这里我使用Channel 在线程之间进行通信,这样可以确保在我们继续ZMQ 代码之前完成打开套接字,它还会在一次写入后使异步线程阻塞。希望您可以对其进行调整以适合您的用例。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-09-12
      • 1970-01-01
      • 1970-01-01
      • 2021-12-24
      • 1970-01-01
      • 2021-05-07
      • 2016-04-05
      • 1970-01-01
      相关资源
      最近更新 更多