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

通过Envoy调用Python GRPC服务出现Stream removed错误排查

问题:gRPC客户端调用报"Stream removed"但服务端处理成功

我在GKE上部署了Envoy代理,用于转发客户端请求到Python gRPC服务。服务端接收Proto消息后会将其发布到Pub/Sub主题,但客户端调用时收到Stream removed错误,尽管服务端已成功处理请求并推送消息到Pub/Sub,客户端却无法收到预期响应。

环境配置与代码

Envoy部署YAML配置

apiVersion: apps/v1
kind: Deployment
metadata:
  name: envoy-deployment
  labels:
    app: envoy
spec:
  replicas: 1
  selector:
    matchLabels:
      app: envoy
  template:
    metadata:
      labels:
        app: envoy
    spec:
      containers:
      - name: envoy
        image: envoyproxy/envoy:v1.22.5
        ports:
        - containerPort: 9901
        livenessProbe:
          httpGet:
            path: /healthz
            port: 9901
          initialDelaySeconds: 60
          timeoutSeconds: 5
          periodSeconds: 10
          failureThreshold: 2
        readinessProbe:
          httpGet:
            path: /healthz
            port: 9901
          initialDelaySeconds: 30
          timeoutSeconds: 5
          periodSeconds: 10
          failureThreshold: 2
        volumeMounts:
        - name: config
          mountPath: /etc/envoy
      volumes:
      - name: config
        configMap:
          name: envoy-conf
---
apiVersion: v1
kind: Service
metadata:
  name: envoy-deployment-service
  annotations:
    cloud.google.com/backend-config: '{"ports": {"9903":"envoy-app-backend-config"}}'
spec:
  ports:
  - protocol: TCP
    port: 9903
    targetPort: 9901
  selector:
    app: envoy
  type: LoadBalancer
  externalTrafficPolicy: Local
---
apiVersion: cloud.google.com/v1
kind: BackendConfig
metadata:
  name: envoy-app-backend-config
spec:
  customRequestHeaders:
    headers:
    - "TE:trailers"
---

apiVersion: networking.k8s.io/v1
kind: Ingress
metadata:
  name: envoy-ingress-prod
  namespace: seshat
  annotations:
    kubernetes.io/ingress.global-static-ip-name: envoy-ingress
    kubernetes.io/ingress.allow-http: "false"
    cert-manager.io/issuer: superset-issuer
    cloud.google.com/backend-config: '{"default": "envoy-app-backend-config"}'

  labels:
    name: envoy-ingress-app
spec:
  tls:
  - hosts:
    - example.com
    secretName: example-tls
  rules:
  - host: example.com
    http:
      paths:
      - path: /*
        pathType: ImplementationSpecific
        backend:
          service:
            name: envoy-deployment-service
            port:
              number: 9903
---
apiVersion: v1
kind: ConfigMap
metadata:
  name: envoy-conf

data:
  envoy.yaml: |
    admin:
      access_log_path: /tmp/admin_access.log
      address:
        socket_address: { address: 127.0.0.1, port_value: 9902 }

    static_resources:
      listeners:
        - name: listener_0
          address:
            socket_address: { address:  0.0.0.0, port_value: 9901 }
          filter_chains:
            - filters:
              - name: envoy.filters.network.http_connection_manager
                typed_config:
                  "@type": type.googleapis.com/envoy.extensions.filters.network.http_connection_manager.v3.HttpConnectionManager
                  codec_type: auto
                  access_log:
                  - name: envoy.access_loggers.file
                    typed_config:
                      "@type": type.googleapis.com/envoy.extensions.access_loggers.file.v3.FileAccessLog
                      path: "/dev/stdout"
                      typed_json_format:
                        "@timestamp": "%START_TIME%"
                        client.address: "%DOWNSTREAM_REMOTE_ADDRESS%"
                        client.local.address: "%DOWNSTREAM_LOCAL_ADDRESS%"
                        envoy.route.name: "%ROUTE_NAME%"
                        envoy.upstream.cluster: "%UPSTREAM_CLUSTER%"
                        host.hostname: "%HOSTNAME%"
                        http.request.body.bytes: "%BYTES_RECEIVED%"
                        http.request.duration: "%DURATION%"
                        http.request.headers.bytes: "%REQUEST_HEADERS_BYTES%"
                        http.request.headers.accept: "%REQ(ACCEPT)%"
                        http.request.headers.authority: "%REQ(:AUTHORITY)%"
                        http.request.headers.te: "%REQ(:TE)%"
                        http.request.headers.id: "%REQ(X-REQUEST-ID)%"
                        http.request.headers.x_forwarded_for: "%REQ(X-FORWARDED-FOR)%"
                        http.request.headers.x_forwarded_proto: "%REQ(X-FORWARDED-PROTO)%"
                        http.request.headers.x_b3_traceid: "%REQ(X-B3-TRACEID)%"
                        http.request.headers.x_b3_parentspanid: "%REQ(X-B3-PARENTSPANID)%"
                        http.request.headers.x_b3_spanid: "%REQ(X-B3-SPANID)%"
                        http.request.headers.x_b3_sampled: "%REQ(X-B3-SAMPLED)%"
                        http.request.method: "%REQ(:METHOD)%"
                        http.response.body.bytes: "%BYTES_SENT%"
                  stat_prefix: ingress_http
                  route_config:
                    name: local_route
                    virtual_hosts:
                      - name: envoy_service
                        domains: ["*"]
                        routes:
                        - match:
                           prefix: "/healthz"
                          direct_response: { status: 200, body: { inline_string: "ok it is working now" } }
                        - match:
                           prefix: "/heal"
                          direct_response: { status: 200, body: { inline_string: "ok heal is working now" } }
                        - match:
                           prefix: "/"
                           #headers:
                           #  - name: te
                           #    exact_match: "trailers"
                          route: {
                            prefix_rewrite: "/",
                            cluster: envoy_service
                          }
                        cors:
                          allow_origin_string_match:
                            - prefix: "*"
                          allow_methods: GET, PUT, DELETE, POST, OPTIONS
                          allow_headers: keep-alive,user-agent,cache-control,content-type,content-transfer-encoding,custom-header-1,x-accept-content-transfer-encoding,x-accept-response-streaming,x-user-agent,x-grpc-web,grpc-timeout
                          max_age: "1728000"
                          expose_headers: custom-header-1,grpc-status,grpc-message
                  http_filters:
                    - name: envoy.filters.http.cors
                      typed_config:
                        "@type": type.googleapis.com/envoy.extensions.filters.http.cors.v3.Cors
                    - name: envoy.filters.http.grpc_web
                      typed_config:
                        "@type": type.googleapis.com/envoy.extensions.filters.http.grpc_web.v3.GrpcWeb
                    - name: envoy.filters.http.router
                      typed_config:
                        "@type": type.googleapis.com/envoy.extensions.filters.http.router.v3.Router
      clusters:
        - name: envoy_service
          connect_timeout: 0.25s
          type: strict_dns
          http2_protocol_options: {}
          lb_policy: round_robin
          load_assignment:
            cluster_name: envoy_service
            endpoints:
              - lb_endpoints:
                - endpoint:
                    address:
                      socket_address:
                        address: seshat-app-server-headless
                        port_value: 8000

Python gRPC服务端代码

server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
master_pb2_grpc.add_EventBusOneofServiceServicer_to_server(
    EventBusServiceServicer(), server
)
server.add_insecure_port("0.0.0.0:8000")
server.start()
print("server started")

def handle_sigterm(*_):
    print("Received shutdown signal")
    all_rpcs_done_event = server.stop(30)
    all_rpcs_done_event.wait(30)
    print("Shut down gracefully")

signal(SIGTERM, handle_sigterm)
server.wait_for_termination()

Python gRPC客户端代码

class ExampleServiceClient(object):
    def __init__(self):
        """Initializer.
           Creates a gRPC channel for connecting to the server.
           Adds the channel to the generated client stub.
        Arguments:
            None.

        Returns:
            None.
        """
       
        self.channel = grpc.secure_channel("domain name", grpc.ssl_channel_credentials(), options=(('grpc.enable_http_proxy', 0),))
        self.stub = master_pb2_grpc.EventBusOneofServiceStub(self.channel)

    def receiveEvent(self, request):
        """Gets a user.
        Arguments:
            name: The resource name of a user.

        Returns:
            None; outputs to the terminal.
        """

        try:
            print(request)
         
            response = self.stub.ReceiveOneofEvent(request)
            print("User fetched.")
            print(response)
        except grpc.RpcError as err:
            print(err)
            print(err.details())  # pylint: disable=no-member
            print("{}, {}".format(err.code().name, err.code().value))  #


if __name__ == "__main__":
    os.environ['GRPC_TRACE'] = 'all'
    os.environ['GRPC_VERBOSITY'] = 'DEBUG'
    if os.environ.get('https_proxy'):
        print("yes proxy present")
        del os.environ['https_proxy']
    if os.environ.get('http_proxy'):
        print("yes proxy present")
        del os.environ['http_proxy']

    for x in range(1,2):
        client = ExampleServiceClient()
        from google.protobuf.json_format import Parse, MessageToJson
        msg = master_pb2.ReceiveOneofEventRequest()
        msg.r.first_name = "a"
        msg.r.last_name = "b"
        msg.r.email = "c"
        client.receiveEvent(msg)

客户端错误信息

<_InactiveRpcError of RPC that terminated with:
    status = StatusCode.UNKNOWN
    details = "Stream removed"
    debug_error_string = "{"created":"@1667812262.198477000","description":"Error received from peer ipv4:IP:443","file":"src/core/lib/surface/call.cc","file_line":967,"grpc_message":"Stream removed","grpc_status":2}"
>
Stream removed
UNKNOWN, (2, 'unknown')

可能的原因与解决方法

1. gRPC-Web过滤器与原生gRPC客户端冲突

原因:Envoy配置中启用了grpc_web过滤器,该过滤器是为浏览器端的gRPC-Web客户端设计的,会修改gRPC的响应格式(比如将HTTP/2 trailers转为响应头),但原生gRPC客户端依赖标准的HTTP/2 trailers来接收响应状态,两者不兼容,导致客户端认为流被中断。

解决方法:如果不需要支持gRPC-Web客户端,直接移除grpc_web过滤器:

http_filters:
  - name: envoy.filters.http.cors
    typed_config:
      "@type": type.googleapis.com/envoy.extensions.filters.http.cors.v3.Cors
  # 移除gRPC-Web过滤器配置
  # - name: envoy.filters.http.grpc_web
  #   typed_config:
  #     "@type": type.googleapis.com/envoy.extensions.filters.http.grpc_web.v3.GrpcWeb
  - name: envoy.filters.http.router
    typed_config:
      "@type": type.googleapis.com/envoy.extensions.filters.http.router.v3.Router

2. HTTP/2 Trailers处理异常

原因:gRPC依赖HTTP/2 trailers传递响应状态和元数据,虽然BackendConfig添加了TE:trailers请求头,但Envoy可能未正确配置以保留或转发响应的trailers,导致客户端无法识别响应完成信号,认为流被提前关闭。

解决方法:在Envoy的HttpConnectionManager配置中添加以下项,确保正确处理gRPC的HTTP/2特性:

typed_config:
  "@type": type.googleapis.com/envoy.extensions.filters.network.http_connection_manager.v3.HttpConnectionManager
  codec_type: auto
  generate_request_id: true
  upgrade_configs:
    - upgrade_type: h2c
  # 保留其他现有配置

3. 连接超时设置过短

原因:Envoy的cluster或GKE负载均衡的空闲/请求超时时间过短,服务端处理Pub/Sub发布的耗时超过了超时时间,导致连接被提前关闭,客户端收到Stream removed错误,但服务端已完成处理。

解决方法:

  • 在Envoy的cluster配置中添加idle_timeout:
    clusters:
      - name: envoy_service
        connect_timeout: 0.25s
        type: strict_dns
        http2_protocol_options: {}
        idle_timeout: 30s  # 设置足够覆盖请求处理的超时时间
        lb_policy: round_robin
        # 保留其他现有配置
    
  • 检查GKE Ingress和LoadBalancer的超时设置,确保其超时时间大于服务端处理请求的最大耗时。

4. 服务端未正确返回响应

原因:Python服务端的ReceiveOneofEvent方法可能未正确构造或返回响应对象,导致Envoy无法获取完整响应,进而关闭流。

解决方法:检查服务端方法实现,确保返回正确的响应对象:

class EventBusServiceServicer(master_pb2_grpc.EventBusOneofServiceServicer):
    def ReceiveOneofEvent(self, request, context):
        # 处理Pub/Sub发布逻辑
        publish_message_to_pubsub(request)
        # 必须返回定义好的响应对象
        return master_pb2.ReceiveOneofEventResponse(
            status="success",
            message="Message published to Pub/Sub"
        )

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 18:15:38