Hacktoberfest 2026:维护者为十月标记出来的 issue,仍然开放、适合新手。 浏览 Hacktoberfest issue

Reading fails when using DirectCodecFactory

未关闭
#3,150 0 条评论 0 个 reaction 已指派 0 人 在 GitHub 查看

还没有人认领这个 Issue。

评估

难度
4/5
预计耗时
3-5 天
新手友好度
45/100
Issue 类型
缺陷
描述清晰度
基本清楚
活跃度
停滞
技术栈
java
领域
data

调研方向

从提供的 reproducer 和 DirectCodecFactory.createDirectCodecFactory 调用开始,然后将其行为与 ParquetReader 使用的 Hadoop codec factory 进行比较。使用附加的 test.parquet 文件,并验证读取是否以 51,000 条记录完成且保留预期的键/值内容。

由索引模型根据 Issue 内容生成。

描述

Type: bug
Describe the bug, including details regarding any error messages, version, and platform.

Hello folks,

I'm currently working on removing the hadoop-common dependency from the runtime in one of my projects.

As part of this process, I need to replace the "Hadoop" codec factory, which relies on certain classes from hadoop-common, with a codec factory that exclusively uses classes from parquet-hadoop. One potential option is DirectCodecFactory, but it comes with a problem.

For demonstration purposes, I’m using a simple key-value Snappy-compressed parquet file containing 51,000 records. Here’s a sample of the data:
Image

When attempting to read this file using DirectCodecFactory, I encounter two issues:

  1. The file fails to read completely. Near the end, it throws an error: Can't read value in column [key] optional binary key (STRING) = 0 at value 49,534 out of 51,000, 9,534 out of 11,000 in currentPage. Repetition level: 0, definition level: 1.
  2. At record number 40,001, the key and value columns get mixed up, with the key column unexpectedly containing a value.

Parquet version: 1.15.0
Hadoop version: 3.4.1

Observations:

  1. With the Hadoop codec factory the file can be read without any issues.

Here are tests demonstrating the issues:

package sandbox.parquet.reader;

import org.apache.hadoop.conf.Configuration;
import org.apache.parquet.bytes.DirectByteBufferAllocator;
import org.apache.parquet.column.ParquetProperties;
import org.apache.parquet.conf.ParquetConfiguration;
import org.apache.parquet.conf.PlainParquetConfiguration;
import org.apache.parquet.hadoop.CodecFactory;
import org.apache.parquet.hadoop.ParquetReader;
import org.apache.parquet.hadoop.api.InitContext;
import org.apache.parquet.hadoop.api.ReadSupport;
import org.apache.parquet.io.InputFile;
import org.apache.parquet.io.LocalInputFile;
import org.apache.parquet.io.api.Binary;
import org.apache.parquet.io.api.Converter;
import org.apache.parquet.io.api.GroupConverter;
import org.apache.parquet.io.api.PrimitiveConverter;
import org.apache.parquet.io.api.RecordMaterializer;
import org.apache.parquet.schema.GroupType;
import org.apache.parquet.schema.MessageType;
import org.apache.parquet.schema.Type;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;

import java.nio.file.Path;
import java.nio.file.Paths;
import java.util.Map;

public class ParquetReaderTest {

    @Test
    void readParquetFileAndVerifyRecordCount_usingDirectCodeFactory() throws Exception {
        Path filePath = Paths.get(ClassLoader.getSystemResource("test.parquet").toURI());
        int recordCount = 0;
        try(ParquetReader<String[]> parquetReader = createReaderWithDirectCodeFactory(filePath)) {
            while (parquetReader.read() != null) {
                recordCount++;
            }
        }

        Assertions.assertEquals(51_000, recordCount);
    }

    @Test
    void readParquetFileAndVerifyContent_usingDirectCodeFactory() throws Exception {
        Path filePath = Paths.get(ClassLoader.getSystemResource("test.parquet").toURI());
        int recordCount = 0;
        try(ParquetReader<String[]> parquetReader = createReaderWithHadoopCodeFactory(filePath)) {
            String[] record;
            while ((record = parquetReader.read()) != null) {
                Assertions.assertEquals("key_" + (recordCount + 1), record[0]);
                Assertions.assertEquals("value_" + (recordCount + 1), record[1]);
                recordCount++;
            }
        }

        Assertions.assertEquals(51_000, recordCount);
    }

    @Test
    void readParquetFileAndVerifyRecordCount_usingHadoopCodeFactory() throws Exception {
        Path filePath = Paths.get(ClassLoader.getSystemResource("test.parquet").toURI());
        int recordCount = 0;
        try(ParquetReader<String[]> parquetReader = createReaderWithHadoopCodeFactory(filePath)) {
            while (parquetReader.read() != null) {
                recordCount++;
            }
        }

        Assertions.assertEquals(51_000, recordCount);
    }

    @Test
    void readParquetFileAndVerifyContent_usingHadoopCodeFactory() throws Exception {
        Path filePath = Paths.get(ClassLoader.getSystemResource("test.parquet").toURI());
        int recordCount = 0;
        try(ParquetReader<String[]> parquetReader = createReaderWithDirectCodeFactory(filePath)) {
            String[] record;
            while ((record = parquetReader.read()) != null) {
                Assertions.assertEquals("key_" + (recordCount + 1), record[0]);
                Assertions.assertEquals("value_" + (recordCount + 1), record[1]);
                recordCount++;
            }
        }

        Assertions.assertEquals(51_000, recordCount);
    }

    private ParquetReader<String[]> createReaderWithDirectCodeFactory(Path file) throws Exception {
        return new ParquetReaderBuilder(new LocalInputFile(file), new PlainParquetConfiguration())
                .withCodecFactory(CodecFactory.createDirectCodecFactory(null,
                        DirectByteBufferAllocator.getInstance(),
                        ParquetProperties.DEFAULT_PAGE_SIZE))
                .build();
    }

    private ParquetReader<String[]> createReaderWithHadoopCodeFactory(Path file) throws Exception {
        // the Hadoop codec factory is created in ParquetReadOptions.Builder#build
        return new ParquetReaderBuilder(new LocalInputFile(file), new PlainParquetConfiguration())
                .build();
    }


    static class ParquetReaderBuilder extends ParquetReader.Builder<String[]> {

        ParquetReaderBuilder(InputFile file, ParquetConfiguration conf) {
            super(file, conf);
        }

        @Override
        protected ReadSupport<String[]> getReadSupport() {
            return new TestReadSupport();
        }
    }

    static class TestReadSupport extends ReadSupport<String[]> {


        @Override
        public ReadContext init(InitContext context) {
            return new ReadContext(context.getFileSchema());
        }

        @Override
        public RecordMaterializer<String[]> prepareForRead(Configuration configuration,
                                                           Map<String, String> keyValueMetaData,
                                                           MessageType fileSchema,
                                                           ReadContext readContext) {
            return new TestRecordMaterializer(fileSchema);
        }

        @Override
        public RecordMaterializer<String[]> prepareForRead(ParquetConfiguration configuration,
                                                           Map<String, String> keyValueMetaData,
                                                           MessageType fileSchema,
                                                           ReadContext readContext) {
            return new TestRecordMaterializer(fileSchema);
        }
    }

    static class TestRecordMaterializer extends RecordMaterializer<String[]> {

        private final TestRootGroupConverter root;

        TestRecordMaterializer(MessageType schema) {
            this.root = new TestRootGroupConverter(schema);
        }

        @Override
        public String[] getCurrentRecord() {
            return root.getCurrentRecord();
        }

        @Override
        public GroupConverter getRootConverter() {
            return root;
        }
    }

    static class TestRootGroupConverter extends GroupConverter {
        private String[] currentRecord;
        private final Converter[] converters;

        TestRootGroupConverter(GroupType schema) {
            converters = new Converter[schema.getFieldCount()];

            for (int i = 0; i < converters.length; i++) {
                final Type type = schema.getType(i);
                if (type.isPrimitive()) {
                    converters[i] = new TestPrimitiveConverter(this, i);
                } else {
                    throw new RuntimeException("Nested records not supported!");
                }
            }
        }

        @Override
        public Converter getConverter(int fieldIndex) {
            return converters[fieldIndex];
        }

        @Override
        public void start() {
            currentRecord = new String[converters.length];
        }

        @Override
        public void end() {
        }

        String[] getCurrentRecord() {
            return currentRecord;
        }
    }

    static class TestPrimitiveConverter extends PrimitiveConverter {

        private final TestRootGroupConverter parent;
        private final int index;

        TestPrimitiveConverter(TestRootGroupConverter parent, int index) {
            this.parent = parent;
            this.index = index;
        }

        @Override
        public void addBinary(Binary value) {
            parent.getCurrentRecord()[index] = value.toStringUsingUTF8();
        }

        @Override
        public void addBoolean(boolean value) {
            throw new UnsupportedOperationException();
        }

        @Override
        public void addDouble(double value) {
            throw new UnsupportedOperationException();
        }

        @Override
        public void addFloat(float value) {
            throw new UnsupportedOperationException();
        }

        @Override
        public void addInt(int value) {
            throw new UnsupportedOperationException();
        }

        @Override
        public void addLong(long value) {
            throw new UnsupportedOperationException();
        }
    }
}

Attaching the parquet file and a demo application:

Component(s)

Core

主要语言
Java
星标
3.1k
派生
1.6k
平均合并
6 天 16 小时
30 天内合并 PR
36

贡献指南

这个仓库没有索引到贡献指南

从这里开始

  1. 先读完整个 Issue,再读项目的贡献指南。
  2. 在 Issue 下留言说明你要接手 —— 这能避免两个人做同样的事。
  3. Fork 仓库,在一个分支上完成修改。
  4. 提交 Pull Request,并在描述里引用这个 Issue 编号。

apache/parquet-java 的其他 Issue

查看 apache/parquet-java 的全部 Issue

相似的 Issue

更多 Java Issue

把新 issue 发到你的邮箱

精选适合新手参与的 GitHub issue 摘要。