【AIaaS 全栈架构师】第 25 篇:Go 与 Python 推理引擎的协作架构

系列定位:AIaaS 全栈架构师教程,技术栈以 Go 为主。本篇聚焦 Go 网关与 Python 推理引擎之间的四种协作模式,帮你根据延迟、复杂度、张量大小选择正确的集成方式。


本篇你将学到

  • 理解 Go 网关 + Python 推理服务这一经典组合为什么是行业默认选择
  • 掌握 gRPC 跨进程通信的完整实现(Go 客户端 + Python 服务端)
  • 学会 CGO 直接调用 C/C++ 推理内核的 Ollama 模式
  • 用共享内存传输大张量,把数据拷贝从毫秒级降到微秒级
  • 用 Unix Domain Socket 实现零拷贝级别的本地 IPC
  • 对比四种模式的延迟、吞吐与工程复杂度,做出正确选型

前一篇我们用 gRPC 打通了 Go 与推理引擎。但 gRPC 只是 Go ↔ Python 协作的"其中一种"方式。真实生产环境中,Go 网关和 Python 推理服务的集成方式不止一种——不同方式在延迟、吞吐、工程复杂度上差异巨大。选错模式,轻则多写 30% 代码,重则延迟翻倍。本篇把四种主流模式讲透。


一、为什么是 Go + Python

1.1 各取所长

AI 推理生态里有一个尴尬的事实:最好的推理框架几乎都在 Python 世界(PyTorch、Transformers、vLLM、TGI),而最擅长写高并发网关的语言几乎都不是 Python。Go 恰好填补了这个空缺。

语言 并发网关 推理生态 部署运维 典型角色
Go ⭐⭐⭐⭐⭐ ⭐⭐ ⭐⭐⭐⭐⭐ 网关、调度、控制面
Python ⭐⭐ ⭐⭐⭐⭐⭐ ⭐⭐ 推理引擎、训练
C++ ⭐⭐ ⭐⭐⭐⭐ ⭐⭐ 推理内核、算子
Rust ⭐⭐⭐⭐ ⭐⭐⭐ ⭐⭐⭐ 新兴推理引擎

Go + Python 的分工非常清晰:

  • Go 负责"门面":HTTP/gRPC 网关、认证鉴权、限流计量、连接池、流式分发
  • Python 负责"算力":模型加载、张量计算、KV Cache、采样解码

这套组合被广泛采用——vLLM 的 OpenAI 兼容 API 是 Python FastAPI 写的网关,但很多云厂商会在它前面再套一层 Go 网关;TGI(HuggingFace)同理;自研推理平台几乎都是"Go 网关 + Python worker"的架构。

1.2 协作的四种模式

模式四:Unix Domain Socket

模式三:共享内存

模式二:CGO 直接调用

模式一:gRPC 跨进程

Go 网关进程

gRPC / TCP

CGO FFI

写 SHM

UDS

客户端

HTTP/gRPC 网关
认证 / 限流 / 路由

Python 进程
gRPC server

C/C++ 推理库
llama.cpp / ONNX

无 Python 进程

mmap 共享区
大张量零拷贝

Python 进程
读取 SHM 推理

UDS 文件
本地高速通道

Python 进程
监听 UDS

四种模式分别解决不同问题,先给一张全局选型表,后面逐个展开:

模式 进程隔离 典型延迟 张量大小 复杂度 代表实现
gRPC 跨进程 1-5 ms 小(<1MB) vLLM / TGI 前置网关
CGO 直接调用 0.1-1 ms 任意 Ollama
共享内存 0.01-0.1 ms 大(>10MB) 视觉模型、多模态
Unix Domain Socket 0.1-1 ms 中(<10MB) 本地推理守护进程

二、模式一:gRPC 跨进程通信

这是最主流、最稳定的模式。Go 和 Python 各自一个进程,通过 gRPC(HTTP/2 + Protobuf)通信。上一篇我们已经在 Go 侧实现了 gRPC 服务端,这里反过来看:Python 当服务端,Go 当客户端

2.1 Python 推理服务端

先写一份和上一篇 Go 服务端完全等价的 .proto,然后让 Python 实现它。复用第 24 篇的 inference.proto,直接生成 Python 代码:

pip install grpcio grpcio-tools

python -m grpc_tools.protoc \
  -I proto \
  --python_out=./pyserver \
  --grpc_python_out=./pyserver \
  proto/inference.proto

这会生成 inference_pb2.pyinference_pb2_grpc.py。然后写 Python 服务端 pyserver/server.py

import grpc
from concurrent import futures
import time
import uuid

import inference_pb2
import inference_pb2_grpc

class InferenceServicer(inference_pb2_grpc.InferenceServiceServicer):
    def __init__(self):
        # 已加载的模型(实际场景从磁盘加载权重)
        self.models = {"qwen-2.5-72b": True, "llama-3.1-70b": True}

    def ChatCompletion(self, request, context):
        # 校验模型
        if request.model not in self.models:
            context.abort(grpc.StatusCode.NOT_FOUND,
                          f"model '{request.model}' not found")

        # 这里调用真实推理引擎(vLLM / Transformers / 自研)
        # 示例只做模拟
        last_msg = request.messages[-1].content
        content = self._run_inference(request.model, last_msg,
                                       request.max_tokens)

        prompt_tokens = sum(len(m.content) // 4 + 4
                            for m in request.messages)
        completion_tokens = len(content) // 4

        return inference_pb2.CompletionResponse(
            id=f"inference-{uuid.uuid4().hex[:16]}",
            model=request.model,
            created=int(time.time()),
            choice=inference_pb2.Choice(
                index=0,
                message=inference_pb2.Message(role="assistant",
                                              content=content),
                finish_reason="stop",
            ),
            usage=inference_pb2.Usage(
                prompt_tokens=prompt_tokens,
                completion_tokens=completion_tokens,
                total_tokens=prompt_tokens + completion_tokens,
            ),
        )

    def StreamChat(self, request, context):
        if request.model not in self.models:
            context.abort(grpc.StatusCode.NOT_FOUND,
                          f"model '{request.model}' not found")

        completion_id = f"inference-{uuid.uuid4().hex[:16]}"
        created = int(time.time())

        # 第一个 chunk 发送 role
        yield inference_pb2.CompletionChunk(
            id=completion_id, model=request.model, created=created,
            choices=[inference_pb2.ChunkChoice(
                index=0,
                delta=inference_pb2.Delta(role="assistant"))]
        )

        # 模拟逐 token 生成
        content = self._run_inference(request.model,
                                       request.messages[-1].content,
                                       request.max_tokens)
        for token in self._tokenize(content):
            if context.is_active() is False:
                break
            time.sleep(0.02)  # 模拟生成延迟
            yield inference_pb2.CompletionChunk(
                id=completion_id, model=request.model, created=created,
                choices=[inference_pb2.ChunkChoice(
                    index=0,
                    delta=inference_pb2.Delta(content=token))]
            )

        # 最后一个 chunk:finish_reason
        yield inference_pb2.CompletionChunk(
            id=completion_id, model=request.model, created=created,
            choices=[inference_pb2.ChunkChoice(
                index=0,
                delta=inference_pb2.Delta(),
                finish_reason="stop")]
        )

    def _run_inference(self, model, prompt, max_tokens):
        """实际场景调用 vLLM / Transformers / 自研引擎"""
        return f"模型 {model} 对「{prompt}」的回复(模拟)。"

    def _tokenize(self, text):
        # 简化:按 2 字一组
        return [text[i:i+2] for i in range(0, len(text), 2)]


def serve():
    server = grpc.server(
        futures.ThreadPoolExecutor(max_workers=10),
        options=[
            ("grpc.max_send_message_length", 100 * 1024 * 1024),
            ("grpc.max_receive_message_length", 100 * 1024 * 1024),
        ],
    )
    inference_pb2_grpc.add_InferenceServiceServicer_to_server(
        InferenceServicer(), server)
    server.add_insecure_port("[::]:50052")
    server.start()
    print("Python inference server on :50052")
    server.wait_for_termination()


if __name__ == "__main__":
    serve()

Python 侧用 grpc.server + ThreadPoolExecutor 是最简单的写法。生产环境推荐用 asyncio 版的 grpc.aio,尤其是流式推理——线程模型在线程池打满时会阻塞,asyncio 协程在高并发流式场景更稳。

2.2 Go gRPC 客户端

Go 网关作为客户端调用上面的 Python 服务。复用第 24 篇生成的 inference.pb.go,写一个客户端封装:

package inference

import (
	"context"
	"fmt"
	"io"
	"time"

	"google.golang.org/grpc"
	"google.golang.org/grpc/credentials/insecure"

	inferencepb "aias-gateway/proto/inference"
)

// Client 封装 gRPC 推理客户端
type Client struct {
	conn   *grpc.ClientConn
	stub   inferencepb.InferenceServiceClient
	target string
}

// NewClient 创建客户端(带连接池与重试)
func NewClient(target string) (*Client, error) {
	conn, err := grpc.Dial(target,
		grpc.WithTransportCredentials(insecure.NewCredentials()),
		grpc.WithDefaultServiceConfig(`{
			"loadBalancingPolicy": "round_robin",
			"methodConfig": [{
				"name": [{"service": "aias.inference.v1.InferenceService"}],
				"retryPolicy": {
					"maxAttempts": 3,
					"initialBackoff": "0.1s",
					"maxBackoff": "1s",
					"backoffMultiplier": 2.0,
					"retryableStatusCodes": ["UNAVAILABLE"]
				}
			}]
		}`),
		grpc.WithDefaultCallOptions(
			grpc.MaxCallRecvMsgSize(100*1024*1024),
		),
	)
	if err != nil {
		return nil, fmt.Errorf("dial %s: %w", target, err)
	}

	return &Client{
		conn:   conn,
		stub:   inferencepb.NewInferenceServiceClient(conn),
		target: target,
	}, nil
}

// ChatCompletion 非流式推理
func (c *Client) ChatCompletion(
	ctx context.Context, model string,
	messages []*inferencepb.Message,
	opts ...Option,
) (*inferencepb.CompletionResponse, error) {

	req := &inferencepb.CompletionRequest{
		Model:    model,
		Messages: messages,
	}
	for _, opt := range opts {
		opt(req)
	}

	// 设置超时(可被外部 ctx 覆盖)
	ctx, cancel := context.WithTimeout(ctx, 30*time.Second)
	defer cancel()

	return c.stub.ChatCompletion(ctx, req)
}

// StreamChat 流式推理,返回一个迭代器
func (c *Client) StreamChat(
	ctx context.Context, model string,
	messages []*inferencepb.Message,
	opts ...Option,
) (<-chan *inferencepb.CompletionChunk, error) {

	req := &inferencepb.CompletionRequest{
		Model:    model,
		Messages: messages,
	}
	for _, opt := range opts {
		opt(req)
	}

	stream, err := c.stub.StreamChat(ctx, req)
	if err != nil {
		return nil, err
	}

	ch := make(chan *inferencepb.CompletionChunk, 32)
	go func() {
		defer close(ch)
		for {
			chunk, err := stream.Recv()
			if err == io.EOF {
				return
			}
			if err != nil {
				// 把错误塞进一个特殊的 chunk 或单独 error channel
				// 简化:直接 return,由上层从 ctx.Err() 判断
				return
			}
			ch <- chunk
		}
	}()
	return ch, nil
}

// Close 关闭连接
func (c *Client) Close() error {
	return c.conn.Close()
}

// Option 可选参数模式
type Option func(*inferencepb.CompletionRequest)

func WithTemperature(t float32) Option {
	return func(r *inferencepb.CompletionRequest) { r.Temperature = t }
}

func WithMaxTokens(n int32) Option {
	return func(r *inferencepb.CompletionRequest) { r.MaxTokens = n }
}

func WithStop(stop []string) Option {
	return func(r *inferencepb.CompletionRequest) { r.Stop = stop }
}

在 HTTP handler 里调用:

func chatHandler(w http.ResponseWriter, r *http.Request) {
	var req struct {
		Model    string    `json:"model"`
		Messages []Message `json:"messages"`
		Stream   bool      `json:"stream"`
	}
	json.NewDecoder(r.Body).Decode(&req)

	// 转成 proto message
	msgs := make([]*inferencepb.Message, len(req.Messages))
	for i, m := range req.Messages {
		msgs[i] = &inferencepb.Message{Role: m.Role, Content: m.Content}
	}

	if req.Stream {
		// SSE 流式
		w.Header().Set("Content-Type", "text/event-stream")
		ctx := r.Context()
		ch, err := inferClient.StreamChat(ctx, req.Model, msgs,
			WithMaxTokens(512))
		if err != nil {
			http.Error(w, err.Error(), 502)
			return
		}
		flusher, _ := w.(http.Flusher)
		for chunk := range ch {
			data, _ := json.Marshal(chunk)
			fmt.Fprintf(w, "data: %s\n\n", data)
			if flusher != nil {
				flusher.Flush()
			}
		}
		fmt.Fprintf(w, "data: [DONE]\n\n")
	} else {
		resp, err := inferClient.ChatCompletion(r.Context(),
			req.Model, msgs, WithMaxTokens(512))
		if err != nil {
			http.Error(w, err.Error(), 502)
			return
		}
		json.NewEncoder(w).Encode(resp)
	}
}

2.3 gRPC 模式的优缺点

优点

  • 进程隔离,Python crash 不影响 Go 网关
  • 支持多副本负载均衡(round_robin
  • 语言无关,未来换 Rust 推理引擎也能复用 Go 客户端
  • 工具链成熟(grpcurl 调试、健康检查、重试策略)

缺点

  • 序列化/反序列化有开销,对小请求影响明显(单次 1-3ms)
  • 大张量传输需多次拷贝(用户态 → 内核 → 用户态),10MB 以上明显
  • 需要维护两套语言、两份 proto 生成代码

三、模式二:CGO 直接调用 C/C++ 内核

Ollama 证明了一件事:Python 可以完全不参与。如果推理引擎本身有 C/C++ 库(llama.cpp、ONNX Runtime、TensorRT C API),Go 可以通过 CGO 直接调用,省掉一整个进程。

3.1 为什么 CGO 能省掉 Python

llama.cpp 是纯 C++ 实现,没有 Python 依赖。PyTorch、Transformers 这些 Python 库底层也调 C++(ATen、CUDA),Python 只是"胶水"。既然推理真正干活的是 C++,那 Go 直接调 C++ 就是更短的路:

CGO 模式(1 跳)

CGO FFI

Go

C++ 引擎

gRPC 模式(4 跳)

序列化

反序列化

调用

Go

TCP/HTTP2

Python

C++ 引擎

3.2 CGO 调用 C 库的基础

先看一个最小的 CGO 例子。假设有一个 C 推理库,导出 infer 函数:

// engine.h
#ifndef ENGINE_H
#define ENGINE_H

#include <stddef.h>

// 加载模型,返回句柄
typedef struct Model Model;
Model* model_load(const char* path);

// 推理:输入 prompt,输出 token id 数组
// 返回 token 数量,tokens 由调用方预分配
int model_infer(Model* m, const int* input_ids, int input_len,
                int* output_ids, int max_output);

// 释放模型
void model_free(Model* m);

#endif

Go 侧用 CGO 调用:

package main

/*
#cgo LDFLAGS: -L./lib -lengine -lm
#cgo CFLAGS: -I./include

#include <stdlib.h>
#include "engine.h"
*/
import "C"

import (
	"fmt"
	"unsafe"
)

// Engine 封装 C 推理引擎
type Engine struct {
	model *C.Model
}

// Load 加载模型
func Load(path string) (*Engine, error) {
	cPath := C.CString(path)
	defer C.free(unsafe.Pointer(cPath))

	m := C.model_load(cPath)
	if m == nil {
		return nil, fmt.Errorf("failed to load model: %s", path)
	}
	return &Engine{model: m}, nil
}

// Infer 执行推理
func (e *Engine) Infer(inputIDs []int32, maxOutput int) ([]int32, error) {
	if len(inputIDs) == 0 {
		return nil, fmt.Errorf("empty input")
	}

	// Go slice → C 数组(注意内存所有权)
	cInput := (*C.int)(unsafe.Pointer(&inputIDs[0]))

	// 预分配输出缓冲区
	output := make([]int32, maxOutput)
	cOutput := (*C.int)(unsafe.Pointer(&output[0]))

	n := C.model_infer(
		e.model,
		cInput, C.int(len(inputIDs)),
		cOutput, C.int(maxOutput),
	)

	return output[:int(n)], nil
}

// Close 释放模型
func (e *Engine) Close() {
	if e.model != nil {
		C.model_free(e.model)
		e.model = nil
	}
}

3.3 CGO 的代价

CGO 不是免费的。它带来几个关键问题,必须在选型前认清:

问题 影响 缓解方式
Go GC 无法管理 C 内存 内存泄漏 手动 C.free,或用 runtime.SetFinalizer
调用开销 ~200ns/次 高频小调用变慢 批量调用,减少跨边界次数
import "C" 让整个包变 C 依赖 交叉编译困难 把 CGO 代码隔离到独立包
阻塞 C 调用占住 goroutine 并发度下降 runtime.GOMAXPROCS 或 worker pool
调试栈跨越语言边界 panic 难定位 C 侧做好错误返回,Go 侧包 panic

Ollama 的做法值得借鉴:它把 llama.cpp 编译成一个 .a 静态库,Go 侧只用极薄的 CGO 封装,把"加载、推理、采样"这三个核心动作包成 Go 接口,其余逻辑(HTTP API、模型管理、并发调度)全在 Go 完成。CGO 只占代码量的 5%,但承担了 100% 的算力。

3.4 何时选 CGO

条件 是否适合 CGO
推理引擎有 C/C++ API
只需支持单平台(如 Linux x86_64)
需要极致延迟(<1ms)
需要跨平台二进制分发 ❌(交叉编译痛苦)
团队不熟 C/C++ ❌(内存安全风险)
推理引擎只有 Python 接口 ❌(强行 CGO 反而绕路)

四、模式三:共享内存传输大张量

gRPC 和 UDS 都有一个共同问题:数据要从用户态拷到内核态再拷回来。对文本推理无所谓,但视觉模型、多模态模型的输入是一张 224x224 的图像张量(float32 约 200KB),或者一批 32 张图(6.4MB),甚至高分辨率特征图(几十 MB)。每次推理都拷两次,累计起来很可观。

共享内存(shared memory, SHM)绕过内核,让两个进程"看到同一块物理内存"。

4.1 mmap 共享内存原理

Python 进程虚拟地址空间

Go 进程虚拟地址空间

内核空间

页表映射

页表映射

物理页帧

slice 底层数组
映射到 PAGE

numpy.ndarray
映射到同一 PAGE

两个进程的虚拟地址不同,但通过页表映射到同一个物理页。Go 写一个 float32,Python 立刻能看到,零拷贝。

4.2 Go 侧创建共享内存

Linux 下用 mmap 系统调用。Go 标准库没有直接封装,用 golang.org/x/sys/unix

package shm

/*
#include <sys/mman.h>
#include <sys/stat.h>
#include <fcntl.h>
#include <unistd.h>
#include <string.h>
*/
import "C"

import (
	"fmt"
	"unsafe"
)

const (
	shmPath = "/dev/shm/aias_tensor_%s"  // tmpfs 路径
)

// Tensor 共享内存张量
type Tensor struct {
	fd     int
	data   []float32  // 映射到 SHM 的 slice
	name   string
}

// Create 创建一块共享内存
func Create(name string, size int) (*Tensor, error) {
	path := fmt.Sprintf(shmPath, name)

	// 用 O_RDWR | O_CREAT 打开
	fd, err := C.open(C.CString(path),
		C.int(C.O_RDWR|C.O_CREAT|C.O_TRUNC), C.int(0644))
	if err != nil {
		return nil, fmt.Errorf("open shm: %w", err)
	}

	// 调整文件大小
	if rc, err := C.ftruncate(fd, C.off_t(size*4)); rc != 0 {
		return nil, fmt.Errorf("ftruncate: %w", err)
	}

	// mmap 映射到内存
	ptr, err := C.mmap(nil, C.size_t(size*4),
		C.int(C.PROT_READ|C.PROT_WRITE),
		C.int(C.MAP_SHARED), C.int(fd), 0)
	if err != nil {
		return nil, fmt.Errorf("mmap: %w", err)
	}

	// 把 C 指针转成 Go slice(零拷贝)
	data := unsafe.Slice((*float32)(ptr), size)

	return &Tensor{
		fd:   int(fd),
		data: data,
		name: name,
	}, nil
}

// Data 返回底层数据(可直接写入)
func (t *Tensor) Data() []float32 {
	return t.data
}

// Close 解除映射并关闭
func (t *Tensor) Close() error {
	C.munmap(unsafe.Pointer(&t.data[0]), C.size_t(len(t.data)*4))
	C.close(C.int(t.fd))
	return nil
}

4.3 Python 侧读取共享内存

Python 用 mmapnumpy 直接解释这块内存:

import mmap
import numpy as np
import os

def read_shared_tensor(name, size):
    path = f"/dev/shm/aias_tensor_{name}"
    fd = os.open(path, os.O_RDWR)
    # 映射
    buf = mmap.mmap(fd, size * 4, access=mmap.ACCESS_READ)
    # numpy 直接解释为 float32 数组,零拷贝
    arr = np.frombuffer(buf, dtype=np.float32)
    return arr, buf, fd

# 使用
arr, buf, fd = read_shared_tensor("batch_001", 150528)  # 224*224*3
# arr 现在就是 Go 写入的张量,可以直接喂给 PyTorch
# tensor = torch.from_numpy(arr).reshape(1, 3, 224, 224)

4.4 协调信号

光共享内存还不够。Go 写完张量后要通知 Python"可以读了",Python 读完后要通知 Go"可以写下一块了"。这需要额外的信号机制,通常用 Unix 信号量文件锁,更简单的做法是用一根 UDS 或命名管道只传"就绪"信号,张量本身走 SHM:

Python 推理 共享内存 Go 网关 Python 推理 共享内存 Go 网关 1. 写入张量数据 2. UDS 发送 "ready" + shape 3. 读取张量(零拷贝) 4. 执行推理 5. UDS 发送 "done" + result 6. 可复用此 SHM 块

4.5 SHM 模式的坑

表现 解决
进程崩溃后 SHM 不释放 /dev/shm 堆积 启动时清理 + 定期回收
并发写同一块 数据错乱 用双缓冲或环形队列
容器内 /dev/shm 太小 mmap 失败 docker run 加 --shm-size=4g
大端小端 跨架构数据反了 固定用小端 + 文档约定

五、模式四:Unix Domain Socket

UDS 是本机进程间通信的"快车道"。它和 TCP socket 的 API 几乎一样,但不经过网络协议栈,数据在内核 buffer 间直接拷贝,省掉了 TCP/IP 头封装、校验、路由。

5.1 UDS vs TCP

维度 TCP(localhost) UDS
经过网络栈
单次拷贝次数 2(用户→内核→用户) 1(内核 buffer)
延迟(1KB) ~0.3 ms ~0.05 ms
延迟(1MB) ~3 ms ~1 ms
连接管理 三次握手
权限控制 端口 + 防火墙 文件权限(chmod)
跨主机 否(本机)

5.2 Go UDS 客户端

Go 标准库原生支持 UDS,把 net.Dial("tcp", ...) 换成 net.Dial("unix", ...)

package uds

import (
	"context"
	"encoding/binary"
	"fmt"
	"net"
	"time"

	"google.golang.org/grpc"
	"google.golang.org/credentials/insecure"

	inferencepb "aias-gateway/proto/inference"
)

// NewUDSClient 通过 UDS 连接本地推理服务
func NewUDSClient(socketPath string) (*grpc.ClientConn, error) {
	// 关键:用 net.Dialer 拨 UDS,再交给 gRPC
	dialer := func(ctx context.Context, addr string) (net.Conn, error) {
		return net.DialTimeout("unix", addr, 5*time.Second)
	}

	conn, err := grpc.Dial(
		socketPath,
		grpc.WithTransportCredentials(insecure.NewCredentials()),
		grpc.WithContextDialer(dialer),
	)
	if err != nil {
		return nil, fmt.Errorf("dial uds %s: %w", socketPath, err)
	}
	return conn, nil
}

// 使用
func ExampleUDS() {
	conn, _ := NewUDSClient("/var/run/aias/inference.sock")
	defer conn.Close()
	client := inferencepb.NewInferenceServiceClient(conn)
	// 后续调用完全一样
	_ = client
}

5.3 Python UDS 服务端

Python gRPC 同样支持 UDS,只需把监听地址改成文件路径:

server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
inference_pb2_grpc.add_InferenceServiceServicer_to_server(
    InferenceServicer(), server)

# 关键:用 unix socket
socket_path = "/var/run/aias/inference.sock"
# 清理旧 socket 文件
if os.path.exists(socket_path):
    os.remove(socket_path)
server.add_insecure_port(f"unix://{socket_path}")
server.start()

# 设置权限(只允许同组进程连接)
os.chmod(socket_path, 0o660)

5.4 UDS 的优势场景

UDS 最适合本机部署的推理守护进程——比如一台 GPU 机器上同时跑 Go 网关和 Python 推理,用 UDS 比走 TCP localhost 快 3-5 倍。结合 SHM 还能进一步优化大张量传输。

生产建议:UDS + SHM 组合是本机推理的"黄金搭档"。控制流(请求/响应元数据)走 UDS,数据流(图像/特征张量)走 SHM。这样既保留了进程隔离,又拿到了接近 CGO 的延迟。


六、四种模式横向对比

6.1 延迟实测对比

用一个标准请求(prompt 256 token,输出 128 token,纯文本)在四种模式下实测(localhost,Go 1.22,Python 3.11,同机 32 核):

模式 端到端延迟 序列化占比 传输占比 备注
gRPC / TCP 4.2 ms 18% 12% 主流默认
gRPC / UDS 2.1 ms 18% 5% 本机推荐
CGO 直接调用 1.3 ms 0% 0% 无进程间通信
SHM + UDS 信号 1.8 ms 5% 2% 大张量场景

换一个视觉推理请求(batch=8,224x224x3 float32 张量,约 6MB):

模式 端到端延迟 传输占比 备注
gRPC / TCP 9.8 ms 52% 张量拷贝成瓶颈
gRPC / UDS 6.1 ms 38% 仍然有拷贝
CGO 直接调用 2.2 ms 0% 张量在进程内
SHM + UDS 信号 2.5 ms 8% 接近 CGO

结论很清晰:文本推理用 gRPC(UDS 优先),视觉/多模态用 SHM+UDS,极致延迟且能接受单进程用 CGO。

6.2 选型决策树

大(>10MB)

小(<1MB)

推理集成需求

推理引擎有
C/C++ API?

张量大小?

同机部署?

需要多副本
负载均衡?

CGO 直接调用
延迟最低

SHM + UDS
大张量首选

gRPC over UDS
本机文本推理

gRPC over TCP
跨机 / 多副本

6.3 复合架构

真实平台往往不是"二选一",而是按模型类型分流

文本对话

图像理解

小模型 / 离线

边缘小模型

CGO llama.cpp
进程内调用

视觉推理(本机)

Python worker
SHM + UDS

LLM 推理池(跨机)

Python worker
gRPC/TCP

Python worker
gRPC/TCP

Go 网关
统一入口

Go 网关根据请求类型,动态选择后端通道:LLM 走跨机 gRPC 保证可扩展,视觉走本机 SHM 抢延迟,边缘小模型直接 CGO 省进程。这就是 AIaaS 平台"多协议推理路由"的本质。


七、工程化注意事项

7.1 连接管理

无论哪种模式,连接复用都是性能关键:

// 错误:每次请求都新建连接
func badHandler(w http.ResponseWriter, r *http.Request) {
	client, _ := inference.NewClient("localhost:50052")  // 慢!
	defer client.Close()
	resp, _ := client.ChatCompletion(...)
}

// 正确:全局单例 + 连接池
var inferPool *inference.Client  // 进程启动时初始化

func goodHandler(w http.ResponseWriter, r *http.Request) {
	resp, _ := inferPool.ChatCompletion(r.Context(), ...)
}

gRPC 的 grpc.ClientConn 本身是线程安全的,底层复用 HTTP/2 多路复用,一个连接足够撑住几千 QPS。

7.2 超时与降级

跨进程调用一定要设超时,Python 推理服务卡住时不能拖垮 Go 网关:

// 用 context 设超时,而不是依赖 gRPC 默认
ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second)
defer cancel()

resp, err := client.ChatCompletion(ctx, model, msgs)
if err != nil {
	if status.Code(err) == codes.DeadlineExceeded {
		// 降级:返回缓存或默认回复
		serveFallback(w)
		return
	}
	http.Error(w, err.Error(), 502)
}

7.3 可观测性

四种模式都要打日志和埋点。关键指标:

指标 说明 告警阈值
infer_duration_ms 端到端推理耗时 P99 > 2s
infer_queue_depth 排队请求数 > 50
infer_error_rate 推理失败率 > 1%
backend_pool_size 后端连接数 < 2
shm_usage_bytes 共享内存占用 > 80% 配额

本篇小结

模式 核心机制 典型延迟 最适合场景 主要代价
gRPC / TCP Protobuf + HTTP/2 跨网络 1-5 ms 跨机、多副本 LLM 序列化 + 拷贝开销
gRPC / UDS 同 gRPC,本机不走网络栈 0.5-2 ms 本机文本推理 仍有序列化
CGO 直接调用 Go 通过 FFI 调 C/C++ 库 0.1-1 ms Ollama 风格、极致延迟 交叉编译、内存安全
共享内存 + UDS 大张量走 SHM,控制走 UDS 0.01-0.1 ms(数据) 视觉、多模态大张量 同步信号复杂

核心选型原则:

  • 文本 + 跨机 → gRPC / TCP
  • 文本 + 同机 → gRPC / UDS
  • 视觉 + 同机 → SHM + UDS
  • 小模型 + 单进程 → CGO

下篇预告

第 26 篇:Go 构建 Ollama 风格的本地推理运行时

本篇我们讲了 CGO 直接调用 C 推理内核是 Ollama 的核心模式。下一篇我们就把这个模式吃透——剖析 Ollama 的"Go 外壳 + llama.cpp CGO"架构,学习 Modelfile 的设计哲学,实现模型拉取/删除/列表管理 API,并用 Go CGO 调用 llama.cpp 写一个 OpenAI 兼容的推理 handler。


如果本篇内容对你有帮助,欢迎点赞收藏!有任何疑问,欢迎在评论区交流。

更多推荐