代码之家  ›  专栏  ›  技术社区  ›  user_1357

Apache Beam/Dataflow:KVCoder损坏要解码的Inputstream

  •  1
  • user_1357  · 技术社区  · 7 年前

    CustomKey , CustomValue 我通过Avro提供了编码器: CustomKeyCoder , CustomValueCoder .

    因为我需要按KV[CustomKey,CustomValue]分组,所以我注册了 KVCoder.of(new CustomKeyCoder, new CustomValueCoder) . 自定义编码器将输入/输出流包装为数据输入/输出流,并使用Avro数据写入器/读取器。

    KV 我明白了 Forbidden IOException when reading from InputStream . 如前所述,解码的关键部分工作正常,当输入流传入解码值时会抛出错误。KVCoder对键和值都使用相同的输入流,我猜键解码会读取整个流。为什么会这样?使用Avro有问题吗?

      //Coder
      override def decode(inputStream: InputStream): CustomValue = {
        val dataInputStream = new DataInputStream(inputStream)
        val id = dataInputStream.readShort
        underlying.decode(dataInputStream)
      }
    
     //Underlying
      override def decode(inputStream: InputStream): CustomValue = {
        val decoder = DecoderFactory.get().binaryDecoder(inputStream, null)
        val record = datumReader.read(null, decoder)
        CustomValue.decode(record)
      }
    
    0 回复  |  直到 7 年前