【问题标题】:gRPC client streaming rpc pipeline error.(write after end ERROR)gRPC 客户端流式传输 rpc 管道错误。(在结束错误后写入)
【发布时间】:2020-01-12 13:32:28
【问题描述】:

我正在研究节点运行时的 gRPC 服务器-客户端编程。

我在客户端流式传输 rpc 中遇到错误。请查看以下 rpc 方法签名。

service RouteGuide{
  rpc DataStreaming(stream File) returns (Stats) {}
}

message Stats{
  string msg=1;
}

message File{
  bytes chk=1;
}

我想将文件从客户端上传到服务器。所以我定义了客户端流 rpc。

问题是文件上传只会第一次成功。

当我尝试上传另一个文件时,我得到一个错误。 在错误结束后写入。

我认为我没有很好地处理流。任何人都可以帮助为什么会发生这种情况?谢谢!

// server.js
"use strict";

const grpc = require("grpc");
const protoLoader = require("@grpc/proto-loader");
const path = require("path");
const fs = require("fs");
const stream = require("stream");

const PROTO_PATH = path.join(__dirname, "proto", "route.proto"); //    path.resolve("proto", "route.proto")
const packageDefinition = protoLoader.loadSync(PROTO_PATH, {
  keepCase: false,
  longs: String,
  enums: String,
  defaults: true,
  oneofs: true
});
const routeguide = grpc.loadPackageDefinition(packageDefinition).routeguide;

const myTransformStream = new stream.Transform({
  objectMode: true,
  transform(data, enc, cb) {
    cb(null, data.chk.toString());
  }
});

function dataStreaming(strm, cb) {
  console.log("server : streaming function");
  stream.pipeline(
    strm,
    myTransformStream,
    fs.createWriteStream("output.txt"),
    err => {
      if (err) {
        console.log(`server side error : ${err}`);
        cb(err);
      } else {
        console.log("server side no error");
        cb(null, "server side finish");
      }
    }
  );
}

function getServer() {
  const server = new grpc.Server();
  server.addService(routeguide.RouteGuide.service, {
    DataStreaming: dataStreaming
  });
  return server;
}

if (require.main === module) {
  const routeServer = getServer();
  routeServer.bind("localhost:3333", grpc.ServerCredentials.createInsecure());
  routeServer.start();
}

=============

// client.js
"use strict";

const grpc = require("grpc");
const protoLoader = require("@grpc/proto-loader");
const path = require("path");
const fs = require("fs");
const stream = require("stream");

const PROTO_PATH = path.join(__dirname, "proto", "route.proto"); //    path.resolve("proto", "route.proto")
const packageDefinition = protoLoader.loadSync(PROTO_PATH, {
  keepCase: false,
  longs: String,
  enums: String,
  defaults: true,
  oneofs: true
});
const routeguide = grpc.loadPackageDefinition(packageDefinition).routeguide;
const client = new routeguide.RouteGuide(
  "localhost:3333",
  grpc.credentials.createInsecure()
);

const MyTransform = new stream.Transform({
  objectMode: true,
  transform(chk, enc, cb) {
    cb(null, { chk: chk });
  }
});

function runDataStreaming() {
  console.log("inside run-data-streaming()");

  const strm = client.dataStreaming((err, ret) => {
    if (err) {
      console.log("client : file transfer failed.");
      console.log(err);
    } else {
      console.log("client : file transfer succeeded.");
    }
  });

  stream.pipeline(fs.createReadStream("test.txt"), MyTransform, strm, err => {
    if (err) {
      console.log(err.message);
    } else {
      console.log("pipeline succeeded");
    }
  });
}

if (require.main === module) {
  runDataStreaming();
}

【问题讨论】:

  • 您能否在问题中包含原始错误日志?

标签: javascript node.js grpc node-streams grpc-node


【解决方案1】:

我解决了我的问题。有一些问题。我会解释的。

  1. server.js 中的全局变量myTransformStream 是问题所在。我将代码移到dataStreaming 函数中。

  2. 我还更改了关闭、销毁选项。 具体来说,我在fs.createReadStream()fs.createWriteStream() 中打开了autocloseemitclose 选项 并且...打开了转换流中的emitcloseautodestroy 选项(myTransformStreammyTransform

【讨论】:

    猜你喜欢
    • 2018-09-23
    • 1970-01-01
    • 1970-01-01
    • 2019-06-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-01-08
    • 2020-06-28
    相关资源
    最近更新 更多