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

参数化函数后ZIO项目消息转发失效的问题排查

ZIO-HTTP + Pulsar 消息转发问题:WebSocket重构后消息无法进入Queue

我正在开发一个基于zio-http的服务器,负责将Pulsar Topic的消息转发到WebSocket。原本运行正常,但发现关闭Socket时会重复创建consumer,于是重构代码将consumer移到WebSocket消息处理器外创建并传入,确保连接和断开时使用同一个consumer。现在消息无法转发,甚至无法进入Queue,我接触ZIO才几天,完全搞不懂问题出在哪。

修改前的Socket代码

package com.example.dashback

import zio.*
import zio.stream.{ZSink, ZStream}
import zhttp.http.*
import zhttp.socket.*
import zhttp.service.ChannelEvent
import zhttp.service.ChannelEvent.UserEvent.HandshakeComplete
import zhttp.service.ChannelEvent.{ChannelUnregistered, UserEventTriggered}
import com.example.dashback.pulsar.PulsarService
import org.apache.pulsar.client.api.{Message, Schema}

def log(s: String) = println(s"\u001b[38;5;226;48;5;19m${s}\u001b[0m")
def zlog(s: String) = ZIO.succeed(log(s))

object WebSocketHandler:
    private def createSocket(pulsar: PulsarService) = {
        val consumerName = s"dashboard-events-${scala.util.Random.nextInt(Int.MaxValue).toHexString}"
        val consumer = pulsar.consumer[String](consumerName, Seq("persistent://example/00000000-0000-4000-8000-000000000001/events"), Schema.STRING)

        val socket = Http.collectZIO[WebSocketChannelEvent] {
            case ChannelEvent(channel, UserEventTriggered(HandshakeComplete)) =>
                for
                    c <- consumer
                    s = ZStream.fromQueue(c.queue)
                    f <- s.foreach { msg =>
                        channel.writeAndFlush(WebSocketFrame.text(msg.getValue))
                    }
                yield s

            case ChannelEvent(channel, ChannelUnregistered) =>
                for
                    _ <- ZIO.log("WebSocket closed; closing consumer")
                    c <- consumer
                    _ <- ZIO.attempt(c.consumer.close())
                yield ()

            case e =>
                for
                    _ <- ZIO.log(e.toString)
                yield ()
        }

        socket.toSocketApp.toResponse
    }

    val app: Http[PulsarService & JwtService, Throwable, Request, Response] = Http.collectZIO[Request] {
        case r @ Method.GET -> !! / "ws" =>
            for
                pulsar <- ZIO.service[PulsarService]
                s <- createSocket(pulsar)
            yield s
    }

修改后的Socket代码

重构后将consumer在WebSocket消息处理器外创建并传入,确保连接和断开时使用同一个consumer:

package com.example.dashback

import zio.*
import zio.stream.{ZSink, ZStream}
import zhttp.http.*
import zhttp.socket.*
import zhttp.service.ChannelEvent
import zhttp.service.ChannelEvent.UserEvent.HandshakeComplete
import zhttp.service.ChannelEvent.{ChannelUnregistered, UserEventTriggered}
import com.example.dashback.pulsar.{PulsarConsumer, PulsarService}
import org.apache.pulsar.client.api.{Consumer, Message, Schema}

def log(s: String) = println(s"\u001b[38;5;226;48;5;19m${s}\u001b[0m")
def zlog(s: String) = ZIO.succeed(log(s))

object WebSocketHandler:
    private def createSocket(pulsar: PulsarService) = {
        val consumerName = s"dashboard-events-${scala.util.Random.nextInt(Int.MaxValue).toHexString}"
        val consumer = pulsar.consumer[String](consumerName, Seq("persistent://example/00000000-0000-4000-8000-000000000001/events"), Schema.STRING)

        def socket(c: PulsarConsumer[String]) = Http.collectZIO[WebSocketChannelEvent] {
            case ChannelEvent(channel, UserEventTriggered(HandshakeComplete)) =>
                for
                    _ <- zlog("We get here")
                    _ <- ZStream.fromQueue(c.queue).foreach { msg =>
                        channel.writeAndFlush(WebSocketFrame.text(msg.getValue))
                    }.fork
                    _ <- channel.writeAndFlush(WebSocketFrame.text("Hello\n"))
                yield ()

            case ChannelEvent(channel, ChannelUnregistered) =>
                for
                    _ <- ZIO.log("WebSocket closed; closing consumer")
                    _ <- ZIO.attempt(c.consumer.close())
                yield ()

            case e =>
                for
                    _ <- ZIO.log(e.toString)
                yield ()
        }

        for
            c <- consumer
            s <- socket(c).toSocketApp.toResponse
        yield s
    }

    val app: Http[PulsarService & JwtService, Throwable, Request, Response] = Http.collectZIO[Request] {
        case r @ Method.GET -> !! / "ws" =>
            for
                pulsar <- ZIO.service[PulsarService]
                s <- createSocket(pulsar)
            yield s
    }

Pulsar Service代码

package com.example.dashback.pulsar

import scala.concurrent.{ExecutionContext, Future}
import scala.util.{Failure, Success}
import scala.jdk.CollectionConverters.*
import zio.*
import zio.stream.ZStream
import org.apache.pulsar.client.api.{Consumer, Message, MessageListener, PulsarClient, Schema, SubscriptionType}

import com.example.dashback.{log, zlog}

case class PulsarConsumer[T](consumer: Consumer[T], queue: Queue[Message[T]])

class PulsarServiceImpl extends PulsarService:
    private val client = PulsarClient.builder()
        .serviceUrl("pulsar://pulsar.example.com:6650")
        .build()

    def consumer[T](subscriptionName: String, topics: Seq[String], schema: Schema[T]): Task[PulsarConsumer[T]] =
        def createConsumer =
            ZIO.attempt {
                client.newConsumer(schema)
                    .topics(topics.asJava)
                    .subscriptionName(subscriptionName)
                    .subscriptionType(SubscriptionType.Exclusive)
                    .subscribe()
            }

        def receive(consumer: Consumer[T], queue: Queue[Message[T]]): Task[Message[T]] =
            ZIO.async { cb =>
                consumer.receiveAsync().thenAccept { msg =>
                    cb {
                        for
                            _ <- zlog("Offering message to queue")
                            _ <- queue.offer(msg)
                            _ <- zlog("Acknowledging message")
                            _ <- ZIO.attempt(consumer.acknowledge(msg))
                            _ <- zlog("Acknowledged")
                            v <- ZIO.attempt(msg)
                        yield v
                    }
                }
            }

        for
            q <- Queue.unbounded[Message[T]]
            c <- createConsumer
            _ <- receive(c, q).forever.fork
        yield PulsarConsumer(c, q)

object PulsarServiceImpl:
    val layer = ZLayer.succeed(new PulsarServiceImpl)

调试情况

修改receive函数添加测试消息,日志里既没出现“Offering message to the queue”,也没看到“Putting test message in the queue”,但Socket能收到手动发送的“Hello”。

def receive(consumer: Consumer[T], queue: Queue[Message[T]]): Task[Unit] =
            for
                _ <- zlog("[Pulsar] Putting test message in the queue")
                _ <- queue.offer(MessageImpl.create(new MessageMetadata, ByteBuffer.wrap("Test message".getBytes), schema, topics.head))

                _ <- zlog("[Pulsar] Attempting to receive a message")
                _ <- ZIO.async[Any, Throwable, Message[T]] { cb =>
                    log("[Pulsar] Calling receiveAsync")
                    consumer.receiveAsync().thenAccept { msg =>
                        log("A message was received")
                        cb {
                            for
                                _ <- zlog("Offering message to queue")
                                _ <- queue.offer(msg)
                                _ <- zlog("Acknowledging message")
                                _ <- ZIO.attempt(consumer.acknowledge(msg))
                                _ <- zlog("Acknowledged")
                                v <- ZIO.attempt(msg)
                            yield v
                        }
                    }
                }
            yield ()

更新:注释WebSocket部分,在run方法中直接创建consumer能正常工作,可接收测试消息和Pulsar Topic消息。

def run = for
        _ <- ZIO.log("Dashback starting")
        c <- (new PulsarServiceImpl).consumer("test", Seq("persistent://example/00000000-0000-4000-8000-000000000001/events"), Schema.STRING)
        _ <- ZStream.fromQueue(c.queue).foreach { msg =>
            zlog(msg.getValue)
        }
    yield ()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 17:01:07