使用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
相关产品推荐
相关产品推荐

