Zum Inhalt springen

GraphQL Subscriptions mit Apollo Server und Client

Praktischer Leitfaden zu GraphQL-Subscriptions für Echtzeit: WebSocket-Setup in Apollo Server, Subscription-Resolver, Client-Integration, Skalierung.

5 Min. Lesezeit
Echtzeitdaten fließen durch eine GraphQL-Subscription-Pipeline

Wenn Queries und Mutations nicht ausreichen

GraphQL-Queries holen Daten. Mutations ändern Daten. Aber keines von beidem sagt dem Client, wenn sich Daten geändert haben. Ohne Subscriptions greift der Client auf Polling zurück — er fragt immer wieder "hat sich etwas geändert?" in festen Intervallen. Polling funktioniert, aber es verschwendet Bandbreite, erhöht die Serverlast und führt zu einer Latenz, die dem Polling-Intervall entspricht.

Subscriptions lösen das, indem sie Updates an den Client schicken, sobald sie eintreten. Eine Chat-Nachricht erscheint sofort. Ein Aktienkurs aktualisiert sich in Echtzeit. Ein Deployment-Status ändert sich, ohne dass der Nutzer die Seite neu lädt.

Diese Anleitung baut ein komplettes Subscription-System mit Apollo Server und Apollo Client auf und behandelt die WebSocket-Infrastruktur, Resolver-Patterns, Authentifizierung auf der Subscription-Verbindung und die Produktionsaspekte, die Tutorials oft überspringen.

Server-Setup: Apollo mit WebSocket-Transport

Apollo Server 4 enthält keinen eingebauten Support für Subscriptions. Man kombiniert es mit der Bibliothek graphql-ws für den WebSocket-Transport und Express (oder Fastify) für den HTTP-Transport.

tstypescript
// src/server.ts
import { ApolloServer } from "@apollo/server";
import { expressMiddleware } from "@apollo/server/express4";
import { createServer } from "http";
import express from "express";
import { WebSocketServer } from "ws";
import { useServer } from "graphql-ws/lib/use/ws";
import { makeExecutableSchema } from "@graphql-tools/schema";
import { typeDefs } from "./schema";
import { resolvers } from "./resolvers";
import { Context, createContext, createWsContext } from "./context";
 
const app = express();
const httpServer = createServer(app);
 
const schema = makeExecutableSchema({ typeDefs, resolvers });
 
// WebSocket server for subscriptions
const wsServer = new WebSocketServer({
  server: httpServer,
  path: "/graphql",
});
 
const serverCleanup = useServer(
  {
    schema,
    context: async (ctx): Promise<Context> => {
      return createWsContext(ctx);
    },
    onConnect: async (ctx) => {
      const token = ctx.connectionParams?.authToken;
      if (!token) {
        throw new Error("Missing authentication token");
      }
      // Validate token — reject connection if invalid
      const user = await validateToken(token as string);
      if (!user) return false; // Reject connection
      return true;
    },
    onDisconnect: async (ctx) => {
      console.log("Client disconnected from subscriptions");
    },
  },
  wsServer
);
 
const server = new ApolloServer<Context>({
  schema,
  plugins: [
    {
      async serverWillStart() {
        return {
          async drainServer() {
            await serverCleanup.dispose();
          },
        };
      },
    },
  ],
});
 
await server.start();
 
app.use(
  "/graphql",
  express.json(),
  expressMiddleware(server, {
    context: createContext,
  })
);
 
httpServer.listen(4000, () => {
  console.log("Server running on http://localhost:4000/graphql");
  console.log("Subscriptions on ws://localhost:4000/graphql");
});

Das duale Transport-Setup ist essenziell. HTTP kümmert sich um Queries und Mutations. WebSocket kümmert sich um Subscriptions. Beide teilen sich dasselbe Schema und denselben Endpoint-Pfad, nutzen aber unterschiedliche Protokolle.

Schema-Design für Subscriptions

graphqlgraphql
# src/schema.graphql
type Message {
  id: ID!
  content: String!
  author: User!
  channel: Channel!
  createdAt: DateTime!
}
 
type Channel {
  id: ID!
  name: String!
  messages(last: Int): [Message!]!
  memberCount: Int!
}
 
type Mutation {
  sendMessage(channelId: ID!, content: String!): Message!
  createChannel(name: String!): Channel!
}
 
type Subscription {
  messageAdded(channelId: ID!): Message!
  channelUpdated: Channel!
  userPresenceChanged(channelId: ID!): PresenceEvent!
}
 
type PresenceEvent {
  userId: ID!
  username: String!
  status: PresenceStatus!
  channelId: ID!
}
 
enum PresenceStatus {
  ONLINE
  OFFLINE
  TYPING
}

Subscription-Felder sollten granular sein. Statt eines einzigen onAnythingChanged-Subscriptions definiert man spezifische Subscriptions für jeden Ereignistyp, für den sich ein Client interessieren könnte. So abonnieren Clients nur die Ereignisse, die sie tatsächlich brauchen.

Resolver: PubSub und Subscription-Logik

Das PubSub-System verbindet Mutations (die Ereignisse erzeugen) mit Subscriptions (die sie konsumieren). Für die Entwicklung funktioniert ein In-Memory-PubSub. Für die Produktion nutzt man ein Redis-basiertes PubSub.

tstypescript
// src/pubsub.ts
import { PubSub } from "graphql-subscriptions";
import { RedisPubSub } from "graphql-redis-subscriptions";
import Redis from "ioredis";
 
// Development: in-memory
const devPubSub = new PubSub();
 
// Production: Redis-backed for multi-instance support
const prodPubSub = new RedisPubSub({
  publisher: new Redis(process.env.REDIS_URL!),
  subscriber: new Redis(process.env.REDIS_URL!),
});
 
export const pubsub =
  process.env.NODE_ENV === "production" ? prodPubSub : devPubSub;
 
// Event name constants
export const EVENTS = {
  MESSAGE_ADDED: "MESSAGE_ADDED",
  CHANNEL_UPDATED: "CHANNEL_UPDATED",
  PRESENCE_CHANGED: "PRESENCE_CHANGED",
} as const;
tstypescript
// src/resolvers/mutation.ts
import { pubsub, EVENTS } from "../pubsub";
 
export const mutationResolvers = {
  Mutation: {
    sendMessage: async (
      _: unknown,
      args: { channelId: string; content: string },
      ctx: Context
    ) => {
      const message = await ctx.db.message.create({
        data: {
          content: args.content,
          channelId: args.channelId,
          authorId: ctx.user.id,
        },
        include: {
          author: true,
          channel: true,
        },
      });
 
      // Publish to subscribers
      await pubsub.publish(EVENTS.MESSAGE_ADDED, {
        messageAdded: message,
        channelId: args.channelId,
      });
 
      return message;
    },
  },
};
tstypescript
// src/resolvers/subscription.ts
import { withFilter } from "graphql-subscriptions";
import { pubsub, EVENTS } from "../pubsub";
 
export const subscriptionResolvers = {
  Subscription: {
    messageAdded: {
      subscribe: withFilter(
        () => pubsub.asyncIterableIterator(EVENTS.MESSAGE_ADDED),
        (payload, variables) => {
          // Only send to subscribers watching this specific channel
          return payload.channelId === variables.channelId;
        }
      ),
    },
 
    channelUpdated: {
      subscribe: () =>
        pubsub.asyncIterableIterator(EVENTS.CHANNEL_UPDATED),
    },
 
    userPresenceChanged: {
      subscribe: withFilter(
        () => pubsub.asyncIterableIterator(EVENTS.PRESENCE_CHANGED),
        (payload, variables) => {
          return payload.userPresenceChanged.channelId === variables.channelId;
        }
      ),
    },
  },
};

Die Funktion withFilter ist entscheidend für die Performance. Ohne sie erhält jeder Abonnent jedes Ereignis, und der Client verwirft die irrelevanten. Mit Filterung sendet der Server Ereignisse nur an die Abonnenten, die sich für den jeweiligen Kanal oder die Entität interessieren.

Client-Integration mit Apollo Client

tstypescript
// src/lib/apollo-client.ts
import {
  ApolloClient,
  InMemoryCache,
  split,
  HttpLink,
} from "@apollo/client";
import { GraphQLWsLink } from "@apollo/client/link/subscriptions";
import { createClient } from "graphql-ws";
import { getMainDefinition } from "@apollo/client/utilities";
 
const httpLink = new HttpLink({
  uri: "/graphql",
  credentials: "include",
});
 
const wsLink = new GraphQLWsLink(
  createClient({
    url: "ws://localhost:4000/graphql",
    connectionParams: () => ({
      authToken: getAuthToken(),
    }),
    retryAttempts: 5,
    shouldRetry: () => true,
    on: {
      connected: () => console.log("WS connected"),
      closed: () => console.log("WS closed"),
      error: (err) => console.error("WS error:", err),
    },
  })
);
 
// Route subscription operations to WebSocket, everything else to HTTP
const splitLink = split(
  ({ query }) => {
    const definition = getMainDefinition(query);
    return (
      definition.kind === "OperationDefinition" &&
      definition.operation === "subscription"
    );
  },
  wsLink,
  httpLink
);
 
export const apolloClient = new ApolloClient({
  link: splitLink,
  cache: new InMemoryCache(),
});

Die Funktion split leitet Operationen an den richtigen Transport weiter — Subscriptions gehen über WebSocket, Queries und Mutations über HTTP. Diese Trennung ist wichtig, weil WebSocket-Verbinden langlaufend und zustandsbehaftet sind, während HTTP-Requests zustandslos und einfacher zu load-balancen sind.

React-Komponenten mit Live-Updates

tsxtsx
// ❌ Bad: Polling for new messages
import { useQuery, gql } from "@apollo/client";
 
const MESSAGES_QUERY = gql`
  query Messages($channelId: ID!) {
    channel(id: $channelId) {
      messages(last: 50) {
        id
        content
        author { name }
        createdAt
      }
    }
  }
`;
 
function ChatBad({ channelId }: { channelId: string }) {
  const { data } = useQuery(MESSAGES_QUERY, {
    variables: { channelId },
    pollInterval: 1000, // Wasteful — 1 request/second
  });
  return <MessageList messages={data?.channel?.messages ?? []} />;
}
tsxtsx
// ✅ Good: Initial query + subscription for live updates
import { useQuery, useSubscription, gql } from "@apollo/client";
 
const MESSAGE_SUBSCRIPTION = gql`
  subscription OnMessageAdded($channelId: ID!) {
    messageAdded(channelId: $channelId) {
      id
      content
      author {
        name
        avatar
      }
      createdAt
    }
  }
`;
 
function ChatChannel({ channelId }: { channelId: string }) {
  const { data, loading } = useQuery(MESSAGES_QUERY, {
    variables: { channelId },
  });
 
  useSubscription(MESSAGE_SUBSCRIPTION, {
    variables: { channelId },
    onData: ({ client, data: subData }) => {
      const newMessage = subData.data?.messageAdded;
      if (!newMessage) return;
 
      // Update the Apollo cache with the new message
      client.cache.modify({
        id: client.cache.identify({
          __typename: "Channel",
          id: channelId,
        }),
        fields: {
          messages(existing = []) {
            const newRef = client.cache.writeFragment({
              data: newMessage,
              fragment: gql`
                fragment NewMessage on Message {
                  id
                  content
                  author { name avatar }
                  createdAt
                }
              `,
            });
            return [...existing, newRef];
          },
        },
      });
    },
  });
 
  if (loading) return <ChatSkeleton />;
 
  return <MessageList messages={data?.channel?.messages ?? []} />;
}

Der onData-Callback aktualisiert den Apollo-Cache manuell, wenn ein Subscription-Ereignis eintrifft. Dieser Ansatz ist zuverlässiger als subscribeToMore für komplexe Cache-Updates, weil man die volle Kontrolle darüber hat, wie die neuen Daten mit den bestehenden fusioniert werden.

Verbindungsmanagement und Heartbeats

WebSocket-Verbindungen sterben lautlos. Mobile Netzwerke brechen Verbindungen ab, ohne Close-Frames zu senden. Load Balancer beenden inaktive Verbindungen. Ohne Heartbeats landen Clients bei toten Verbindungen, die nie wieder Updates erhalten.

tstypescript
// Server-side: graphql-ws handles heartbeats automatically via ping/pong
// Client-side: configure keepalive and reconnection
const wsClient = createClient({
  url: "ws://localhost:4000/graphql",
  connectionParams: () => ({
    authToken: getAuthToken(),
  }),
  keepAlive: 10000, // Send ping every 10 seconds
  retryAttempts: Infinity, // Always retry
  retryWait: async (retries: number) => {
    // Exponential backoff with jitter
    const baseDelay = Math.min(1000 * 2 ** retries, 30000);
    const jitter = Math.random() * 1000;
    await new Promise((resolve) =>
      setTimeout(resolve, baseDelay + jitter)
    );
  },
});

Exponentielles Backoff mit Jitter verhindert Thundering-Herd-Probleme, wenn ein Server neu startet und alle Clients gleichzeitig reconnecten wollen. Der Jitter verteilt die Reconnect-Versuche über ein Zeitfenster, statt sie auf denselben Moment zu konzentrieren.

Wichtige Erkenntnisse

GraphQL-Subscriptions verwandeln deine API von Request-Response zu Event-Driven. Die Architektur erfordert ein duales Transport-Setup — HTTP für Queries und Mutations, WebSocket für Subscriptions — mit geteiltem Schema und Authentifizierung.

Serverseitiges Filtern mit withFilter ist der Performance-Schlüssel. Ohne es erhält jeder verbundene Client jedes Ereignis, was nicht skaliert. Ein Redis-basiertes PubSub ermöglicht Multi-Instance-Deployments, indem es Ereignisse über Server-Prozesse hinweg teilt.

Auf Client-Seite kombiniert man initiale Queries mit Subscriptions für die beste Nutzererfahrung: den aktuellen Zustand sofort laden, dann Live-Updates schrittweise anwenden. Konfiguriere immer Keepalive und exponentielles Backoff für die WebSocket-Verbindung — stille Verbindungsabbrüche sind das häufigste Produktionsproblem bei Subscriptions.

Wilfredo Rujel

Wilfredo Rujel

Full-Stack-Softwareentwickler

Diesen Beitrag teilenX