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:
- Your
EventObservableis a no-op: ThesubscribeActualmethod just returns without any logic to emit events, so no observer will ever receive data from it. - 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. - Unconnected GraphQL subscription: Your
EventPublisheris a separate, unlinked component, so it can't receive events from your StreamListener. - Manual bean creation: You're directly instantiating
SubscriptionGraphQlUtilitiesinstead 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
Observableis a passive stream that only emits data when subscribed to (based on itssubscribeActuallogic). ASubjectis 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), useFlowableinstead of rawObservableto handle backpressure (preventing subscribers from being overwhelmed with events).
内容的提问来源于stack exchange,提问作者Gaithri Haridas

