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

使用Dart和gRPC无法接收Salesforce Pub/Sub API事件

Dart gRPC实现Salesforce Pub/Sub API无法接收事件的排查方案

我按照Salesforce官方快速入门指南用Python成功实现了Pub/Sub API的事件订阅与监听,但用Dart结合gRPC实现时遇到问题。尽管逻辑和Python版本一致,Dart代码运行无报错、终端保持活跃(看似订阅成功),但完全接收不到事件——而Python脚本能实时收到我发布的事件。已确认所用Session ID具备订阅目标主题的权限。

Dart实现代码

import 'package:grpc/grpc.dart';
import 'package:pubsub_testing/src/generated/salesforceProtoFile.pbgrpc.dart' as pb_grpc;
import "package:pubsub_testing/src/generated/salesforceProtoFile.pb.dart" as pb2;

// Session ID, instance URL, and tenant ID are provided.

final authMetadata = CallOptions(metadata: {
  'accesstoken': sessionId,
  'instanceurl': instanceUrl,
  'tenantid': tenantId,
});

final channel = ClientChannel(
  'api.pubsub.salesforce.com',
  port: 7443,
);

final stub = pb_grpc.PubSubClient(channel);

Stream<pb2.FetchRequest> fetchReqStream(String topic) async* {
  while (true) {
    yield pb2.FetchRequest(
      topicName: topic,
      replayPreset: pb2.ReplayPreset.LATEST,
      numRequested: 100,
    );
  }
}

Future<void> subscribe(String mySubTopic) async {
  print('Subscribing to $mySubTopic');
  try {
    final substream =
        stub.subscribe(fetchReqStream(mySubTopic), options: authMetadata);
    print("substream: $substream");

    await for (var event in substream) {
      print("Got an event!\n");
      if (event.events.isNotEmpty) {
        print("Number of events received: ${event.events.length}");
        var payloadbytes = event.events[0].event.payload;
        var schemaid = event.events[0].event.schemaId;

        var schema = await stub.getSchema(pb2.SchemaRequest(schemaId: schemaid),
            options: authMetadata);

        print("Got an event!\n");
      } else {
        print("[${DateTime.now()}] The subscription is active.");
      }
    }
  } catch (e) {
    print('An error occurred during subscription: $e');
  }
}

运行状态:无报错,终端持续运行但无事件输出,也没有"订阅活跃"的日志。

可正常运行的Python实现代码

from __future__ import print_function
import threading
import io
import pubsub_api_pb2 as pb2
import pubsub_api_pb2_grpc as pb2_grpc
import time
import json

semaphore = threading.Semaphore(1)

latest_replay_id = None

with grpc.secure_channel('api.pubsub.salesforce.com:7443', grpc.ssl_channel_credentials(None)) as channel:
    # Store Auth 
    sessionid = ''
    instanceurl = ''
    tenantid = ''
    authmetadata = (('accesstoken', sessionid),
    ('instanceurl', instanceurl),
    ('tenantid', tenantid))
    
    # Generate Stub
    stub = pb2_grpc.PubSubStub(channel)

    # Subscribe to event channel
    def fetchReqStream(topic):
        while True:
            semaphore.acquire()
            yield pb2.FetchRequest(
                topic_name = topic,
                replay_preset = pb2.ReplayPreset.LATEST,
                num_requested = 1)
    
    # Decode event message payload
    def decode(schema, payload):
        schema = avro.schema.parse(schema)
        buf = io.BytesIO(payload)
        decoder = avro.io.BinaryDecoder(buf)
        reader = avro.io.DatumReader(schema)
        ret = reader.read(decoder)
        return ret

    # Make the subscribe call
    mysubtopic = "/event/RS_L__ConversationEvent__e"
    print('Subscribing to ' + mysubtopic)
    substream = stub.Subscribe(fetchReqStream(mysubtopic),
            metadata=authmetadata   )
    for event in substream:
        if event.events:
            semaphore.release()
            print("Number of events received: ", len(event.events))
            payloadbytes = event.events[0].event.payload
            schemaid = event.events[0].event.schema_id
            schema = stub.GetSchema(
                    pb2.SchemaRequest(schema_id=schemaid),
                    metadata=authmetadata).schema_json
            decoded = decode(schema, payloadbytes)
            print("Got an event!", json.dumps(decoded), "\n")
        else:
            print("[", time.strftime('%b %d, %Y %l:%M%p %Z'),
            "] The subscription is active.")

排查思路与解决方案

1. 修复gRPC通道的TLS加密配置

Python代码使用了grpc.secure_channel建立加密连接,但Dart的ClientChannel默认是明文通道,而Salesforce Pub/Sub API强制要求TLS加密通信。修改Dart通道初始化代码:

final channel = ClientChannel(
  'api.pubsub.salesforce.com',
  port: 7443,
  options: ChannelOptions(
    credentials: ChannelCredentials.secure(),
  ),
);

这是最可能导致事件无法接收的核心原因——明文连接无法与Salesforce服务正常交互。

2. 修正FetchRequest的流控制逻辑

Salesforce Pub/Sub API要求必须在上一个FetchRequest的响应返回后,才能发送下一个请求。Python代码通过信号量semaphore实现了这个逻辑,但Dart代码是无限循环持续发送请求,会触发服务端限流,导致事件无法推送。

修改Dart的流生成逻辑,改为收到响应后再发送下一个请求:

// 用控制器控制请求发送节奏
Stream<pb2.FetchRequest> fetchReqStream(String topic, StreamController<void> trigger) async* {
  // 发送第一个请求
  yield pb2.FetchRequest(
    topicName: topic,
    replayPreset: pb2.ReplayPreset.LATEST,
    numRequested: 100,
  );
  // 等待触发信号,发送下一个请求
  await for (var _ in trigger.stream) {
    yield pb2.FetchRequest(
      topicName: topic,
      replayPreset: pb2.ReplayPreset.LATEST,
      numRequested: 100,
    );
  }
}

// 订阅函数中更新逻辑
Future<void> subscribe(String mySubTopic) async {
  print('Subscribing to $mySubTopic');
  try {
    final triggerController = StreamController<void>();
    final substream = stub.subscribe(fetchReqStream(mySubTopic, triggerController), options: authMetadata);
    
    await for (var event in substream) {
      if (event.events.isNotEmpty) {
        print("Number of events received: ${event.events.length}");
        // 处理事件逻辑...
        
        // 发送下一个请求的触发信号
        triggerController.add(null);
      } else {
        print("[${DateTime.now()}] The subscription is active.");
      }
    }
  } catch (e) {
    print('An error occurred during subscription: $e');
  }
}

3. 验证Proto文件生成正确性

确保Dart的proto文件是从Salesforce官方Pub/Sub API proto正确生成的,检查字段名映射是否一致:

  • Python的topic_name对应Dart的topicName
  • Python的replay_preset对应Dart的replayPreset
  • Python的num_requested对应Dart的numRequested
    字段名不匹配会导致请求参数无法被服务端解析。

4. 开启gRPC日志排查

添加日志包装器查看请求响应细节:

import 'package:grpc/src/client/logging.dart';

// 用日志包装通道
final stub = pb_grpc.PubSubClient(LoggingClientChannel(channel, LoggingLevel.all));

通过日志可以确认请求是否发送、响应是否接收,以及是否有隐藏的错误信息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 21:50:00