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

Spring Boot GraphQL Java订阅与StreamListener:Observable创建问题求助

Hey there! Let's break down what's going wrong with your code and fix up your understanding of Reactive Java + Spring Cloud Stream here.

Core Issues in Your Current Code

First, let's highlight the main problems that are preventing your Observable from working as expected:

  1. Your EventObservable is a no-op: The subscribeActual method just returns without any logic to emit events, so no observer will ever receive data from it.
  2. Incorrect event emission: Calling eventObservable.getObservable().just(event) creates a new Observable instance (instead of sending the event to your existing Observable) and does nothing with it—this line has zero effect.
  3. Unconnected GraphQL subscription: Your EventPublisher is a separate, unlinked component, so it can't receive events from your StreamListener.
  4. Manual bean creation: You're directly instantiating SubscriptionGraphQlUtilities instead of letting Spring manage it, which means its dependencies won't be injected properly.

Fixed Code Implementation

Let's rewrite the key components to fix these issues:

1. Refactor EventObservable with a Subject

A Subject acts as both an Observable and an Observer—perfect for manually pushing events from your StreamListener to subscribers. We'll use PublishSubject (multicasts events to all active subscribers) and make it thread-safe with toSerialized().

@Component
public class EventObservable {
    // Thread-safe Subject to multicast events
    private final Subject<Event> eventSubject = PublishSubject.<Event>create().toSerialized();

    // Expose only the Observable interface to prevent external misuse of Subject methods
    public Observable<Event> getObservable() {
        return eventSubject;
    }

    // Method to push events into the stream
    public void publishEvent(Event event) {
        eventSubject.onNext(event);
    }
}

2. Update EventConsumer to Publish Events

Instead of creating a useless Observable with just(), call our new publishEvent method to send events into the Subject:

@EnableBinding(ConsumerChannels.class)
public class EventConsumer {
    @Autowired
    private EventObservable eventObservable;

    @StreamListener(ConsumerChannels.EVENT_NOTIFICATION_CHANNEL)
    public void fetchEvent(Message<?> message){
        // Note: In production, you'd parse the Event from the incoming message
        Event event = new Event("123123","google.io","This is an event","Singapore");
        // Push the event into our reactive stream
        eventObservable.publishEvent(event);
    }
}

3. Link GraphQL Subscription to EventObservable

Inject EventObservable into SubscriptionGraphQlUtilities and use it to feed the GraphQL subscription data fetcher. We'll convert the Observable to a Flowable (a reactive stream that supports backpressure) to match GraphQL's Publisher requirement.

@Component
public class SubscriptionGraphQlUtilities {
    private Logger logger = LoggerFactory.getLogger(SubscriptionGraphQlUtilities.class);
    @Value("classpath:schemas.graphqls")
    private Resource schemaResource;
    private final EventObservable eventObservable;
    private final GraphQLSchema graphQLSchema;

    // Constructor injection for EventObservable
    public SubscriptionGraphQlUtilities(EventObservable eventObservable){
        this.eventObservable = eventObservable;
        graphQLSchema = buildSchema();
    }

    private GraphQLSchema buildSchema() {
        Reader streamReader = loadSchemaFile("schemas.graphqls");
        TypeDefinitionRegistry typeRegistry = new SchemaParser().parse(streamReader);
        RuntimeWiring wiring = RuntimeWiring.newRuntimeWiring()
                .type(newTypeWiring("subscription")
                        .dataFetcher("subscribeEvent", getEventsDataFetcher())
                )
                .build();
        return new SchemaGenerator().makeExecutableSchema(typeRegistry, wiring);
    }

    private DataFetcher getEventsDataFetcher(){
        return environment -> {
            // Map Event objects to match your GraphQL subscription fields
            return eventObservable.getObservable()
                    .map(event -> Map.of(
                        "id", event.getId(),
                        "name", event.getName(),
                        "description", event.getDescription(),
                        "location", event.getLocation()
                    ))
                    .toFlowable(BackpressureStrategy.BUFFER); // Convert to Publisher with backpressure support
        };
    }

    // Keep your existing loadSchemaFile and getGraphQLSchema methods
    @SuppressWarnings("SameParameterValue")
    private Reader loadSchemaFile(String name) {
        InputStream stream = getClass().getClassLoader().getResourceAsStream(name);
        return new InputStreamReader(stream);
    }

    public GraphQLSchema getGraphQLSchema() {
        return graphQLSchema;
    }
}

4. Fix SubscriptionConsumer Dependency Injection

Let Spring inject SubscriptionGraphQlUtilities instead of instantiating it manually, and clean up the subscription logic:

@EnableBinding(ConsumerChannels.class)
@Component
public class SubscriptionConsumer {
    private Logger logger = LoggerFactory.getLogger(SubscriptionConsumer.class);
    private final AtomicReference<Subscription> subscriptionRef = new AtomicReference<>();
    private final SubscriptionGraphQlUtilities graphQlUtilities;

    // Constructor injection for SubscriptionGraphQlUtilities
    public SubscriptionConsumer(SubscriptionGraphQlUtilities graphQlUtilities) {
        this.graphQlUtilities = graphQlUtilities;
    }

    Instrumentation instrumentation = new ChainedInstrumentation(
            Collections.singletonList(new TracingInstrumentation())
    );

    @StreamListener(ConsumerChannels.QUERY_INPUT)
    public void fetchSubscription(Message<?> message){
        logger.info("Received subscription query:{}", message.toString());
        try {
            String subscriptionQuery = "subscription {\n subscribeEvent{\n\t\tid\n name\n\t\tdescription\n\t\tlocation\n\t}\n}";
            GraphQL graphQL = GraphQL
                    .newGraphQL(graphQlUtilities.getGraphQLSchema())
                    .instrumentation(instrumentation)
                    .build();
            ExecutionResult executionResult = graphQL.execute(subscriptionQuery);
            Publisher<ExecutionResult> eventStream = executionResult.getData();
            
            eventStream.subscribe(new Subscriber<ExecutionResult>() {
                @Override
                public void onSubscribe(Subscription subscription) {
                    logger.info("Successfully subscribed");
                    subscriptionRef.set(subscription);
                    subscription.request(1); // Request first event
                }

                @Override
                public void onNext(ExecutionResult executionResult) {
                    logger.info("Sending event updates");
                    Object result = executionResult.getData();
                    logger.info("Received event: {}", result);
                    subscriptionRef.get().request(1); // Request next event
                }

                @Override
                public void onError(Throwable throwable) {
                    logger.error("Subscription failed", throwable);
                    subscriptionRef.set(null);
                }

                @Override
                public void onComplete() {
                    logger.info("Subscription completed");
                    subscriptionRef.set(null);
                }
            });
        } catch (Exception e) {
            logger.error("Error handling subscription request", e);
        }
    }
}

Key Concept Clarifications

  • Observable vs Subject: An Observable is a passive stream that only emits data when subscribed to (based on its subscribeActual logic). A Subject is an active stream that lets you manually push events (onNext()) to all subscribers—ideal for bridging imperative code (like StreamListener) to reactive code.
  • Dependency Injection: Never manually instantiate Spring-managed beans (like SubscriptionGraphQlUtilities)—use constructor injection so Spring handles wiring dependencies for you.
  • Backpressure: When working with Publisher (required for GraphQL subscriptions), use Flowable instead of raw Observable to handle backpressure (preventing subscribers from being overwhelmed with events).

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:36:22