【问题标题】:Convert Avro to Parquet in NiFi在 NiFi 中将 Avro 转换为 Parquet
【发布时间】:2019-02-19 08:51:54
【问题描述】:

我想在 NiFi 中将 Avro 文件转换为 Parquet。我知道可以通过 ConvertAvroToORC 处理器转换为 ORC,但我没有找到转换为 Parquet 的解决方案。

我正在通过 ConvertRecord(JsonTreeReader 和 AvroRecordSetWriter)处理器将 JSON 转换为 Avro。之后,我想将 Avro 有效负载转换为 Parquet,然后再将其放入 S3 存储桶中。我不想将它存储在 HDFS 中,因此 PutParquet 处理器似乎不适用。

我需要一个处理器,例如:ConvertAvroToParquet

【问题讨论】:

    标签: apache-nifi


    【解决方案1】:

    @Martin,您可以使用我最近在Nifi 中贡献的非常方便的处理器ConvertAvroToParquet。它应该在最新版本中可用。

    此处理器的用途与您正在寻找的完全相似。有关此处理器的更多详细信息及其创建原因:Nifi-5706

    代码Link

    【讨论】:

      【解决方案2】:

      实际上可以使用 PutParquet 处理器。

      以下描述来自 nifi-1.8 中的工作流程。

      将以下库放入文件夹,例如home/nifi/s3libs/:

      • aws-java-sdk-1.11.455.jar(+ 第三方库)
      • hadoop-aws-3.0.0.jar

      创建一个 xml 文件,例如/home/nifi/s3conf/core-site.xml。可能需要一些额外的调整,为您的区域使用正确的端点。

      <configuration>
          <property>
              <name>fs.defaultFS</name>
              <value>s3a://BUCKET_NAME</value>
          </property>
          <property>
              <name>fs.s3a.access.key</name>
              <value>ACCESS-KEY</value>
          </property>
          <property>
              <name>fs.s3a.secret.key</name>
              <value>SECRET-KEY</value>
          </property>
          <property>
              <name>fs.AbstractFileSystem.s3a.imp</name>
              <value>org.apache.hadoop.fs.s3a.S3A</value>
          </property>
          <property>
              <name>fs.s3a.multipart.size</name>
              <value>104857600</value>
              <description>Parser could not handle 100M. replacing with bytes. Maybe not needed after testing</description>
          </property>
          <property>
              <name>fs.s3a.endpoint</name>
              <value>s3.eu-central-1.amazonaws.com</value> 
              <description>Frankfurt</description>
          </property>
          <property>
              <name>fs.s3a.fast.upload.active.blocks</name>
              <value>4</value>
              <description>
          Maximum Number of blocks a single output stream can have
          active (uploading, or queued to the central FileSystem
          instance's pool of queued operations.
      
          This stops a single stream overloading the shared thread pool.
              </description>
          </property>
          <property>
              <name>fs.s3a.threads.max</name>
              <value>10</value>
              <description>The total number of threads available in the filesystem for data
          uploads *or any other queued filesystem operation*.</description>
          </property>
      
          <property>
              <name>fs.s3a.max.total.tasks</name>
              <value>5</value>
              <description>The number of operations which can be queued for execution</description>
          </property>
      
          <property>
              <name>fs.s3a.threads.keepalivetime</name>
              <value>60</value>
              <description>Number of seconds a thread can be idle before being terminated.</description>
          </property>
          <property>
              <name>fs.s3a.connection.maximum</name>
              <value>15</value>
          </property>
      </configuration>
      

      用法

      创建一个PutParquet 处理器。在属性下设置

      • Hadoop 配置资源:/home/nifi/s3conf/core-site.xml
      • 其他类路径资源:/home/nifi/s3libs
      • 目录:s3a://BUCKET_NAME/folder/(EL 可用)
      • 压缩类型:使用 NONE、SNAPPY 测试
      • 删除 CRC:真

      流文件必须包含 filename 属性 - 没有花哨的字符或斜线。

      【讨论】:

      • 这是在 Nifi-5706 下捕获的
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-06-22
      • 2020-10-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多