如何在Apollo Federation网关GraphQL中跨子图调用Mutation
跨服务联动实现:创建Post后更新用户的Posts数组
现有User和Posts两个独立子图服务,分别使用不同数据库。需求是创建Post后,自动将新生成的Post ID添加到对应用户的posts数组中。其中createPost Mutation属于Post服务,addUserPost Mutation属于User服务。
现有代码片段
Post服务的createPost Mutation
const Posts = require("../models/Posts"); const resolvers = { Mutation: { createPost: async (parent, args, context) => { const { input } = args; const newPost = await Posts.create({ input }); return newPost; } } }; module.exports = { resolvers };
User服务的addUserPost Mutation
const User = require("../models/User"); const resolvers = { Mutation: { addUserPost: async (parent, args, context) => { const { newPostId } = args; const { email } = context; const updateUser = await User.findOneAndUpdate( { email }, { $push: { posts: newPostId } } ); return updateUser; } } }; module.exports = { resolvers };
Apollo联邦网关配置
const { ApolloServer } = require('apollo-server'); const { ApolloGateway, RemoteGraphQLDataSource } = require('@apollo/gateway'); class AuthenticatedDataSource extends RemoteGraphQLDataSource { willSendRequest({ request, context }) { const headers = context?.req?.headers || {}; request.http.headers.set('Authorization', headers.authorization || ''); } } const gateway = new ApolloGateway({ serviceList: [ { name: 'user', url: 'http://localhost:5000' }, { name: 'post', url: 'http://localhost:4000' }, ], buildService: ({ url }) => new AuthenticatedDataSource({ url }), }); const server = new ApolloServer({ gateway, subscriptions: false, context: ({ req }) => ({ req }) }); server.listen(8000, () => { console.log(`Gateway Server is running on port 8000`); });
实现方案
方案一:网关层定义组合Mutation
在网关层封装一个顶层Mutation,依次调用两个服务的接口,保证操作连贯性:
修改网关代码,添加自定义组合Mutation的resolver:
const { ApolloServer } = require('apollo-server'); const { ApolloGateway, RemoteGraphQLDataSource } = require('@apollo/gateway'); class AuthenticatedDataSource extends RemoteGraphQLDataSource { willSendRequest({ request, context }) { const headers = context?.req?.headers || {}; request.http.headers.set('Authorization', headers.authorization || ''); } } const gateway = new ApolloGateway({ serviceList: [ { name: 'user', url: 'http://localhost:5000' }, { name: 'post', url: 'http://localhost:4000' }, ], buildService: ({ url }) => new AuthenticatedDataSource({ url }), }); const server = new ApolloServer({ gateway, subscriptions: false, context: ({ req }) => ({ req }), resolvers: { Mutation: { createPostAndLinkToUser: async (_, args, { req, gateway }) => { // 1. 调用Post服务创建帖子 const postResult = await gateway.execute({ document: ` mutation CreatePost($input: PostInput!) { createPost(input: $input) { id title content } } `, variables: { input: args.input }, context: { req }, }); if (postResult.errors) throw postResult.errors[0]; const newPost = postResult.data.createPost; // 2. 调用User服务关联帖子ID到用户 const userResult = await gateway.execute({ document: ` mutation AddUserPost($newPostId: ID!) { addUserPost(newPostId: $newPostId) { id email posts } } `, variables: { newPostId: newPost.id }, context: { req }, }); if (userResult.errors) throw userResult.errors[0]; const updatedUser = userResult.data.addUserPost; return { post: newPost, user: updatedUser }; } } } }); server.listen(8000, () => { console.log(`Gateway Server is running on port 8000`); });
方案二:Post服务直接调用User服务接口
在Post服务的createPost resolver中,创建帖子后直接调用User服务的addUserPost接口:
- 先安装请求工具(比如axios):
npm install axios
- 修改Post服务的resolver:
const Posts = require("../models/Posts"); const axios = require('axios'); const resolvers = { Mutation: { createPost: async (parent, args, context) => { const { input } = args; const newPost = await Posts.create(input); // 注意:根据你的模型定义,可能需要调整为{ ...input } try { // 调用User服务关联帖子ID await axios.post('http://localhost:5000/graphql', { query: ` mutation AddUserPost($newPostId: ID!) { addUserPost(newPostId: $newPostId) { id } } `, variables: { newPostId: newPost.id }, headers: { Authorization: context.req.headers.authorization || '' } }); } catch (err) { // 关联失败时回滚帖子创建 await Posts.findByIdAndDelete(newPost.id); throw new Error('创建帖子后关联用户失败,请重试'); } return newPost; } } }; module.exports = { resolvers };
关键注意事项
- 数据一致性:跨服务操作无法依赖数据库事务,需添加补偿机制(如方案二中的回滚),或使用MQ等事件驱动架构解耦服务。
- 权限校验:确保调用User服务时携带正确的身份凭证,保证
context.email能正确获取当前用户。 - 错误处理:必须处理单个服务调用失败的情况,避免出现帖子已创建但未关联到用户的不一致状态。
内容的提问来源于stack exchange,提问作者Harshil Sharma
相关产品推荐
相关产品推荐

