<dependency> <groupId>org.apache.iceberg</groupId> <artifactId>api</artifactId> <version>0.11.0</version> </dependency> import org.apache.iceberg.*; import org.apache.iceberg.data.*; import org.apache.iceberg.expressions.*; import org.apache.iceberg.hadoop.*; import org.apache.iceberg.types.*; import org.apache.iceberg.parquet.*; Schema schema = new Schema( Types.NestedField.required(1, "id", Types.IntegerType.get()), Types.NestedField.required(2, "name", Types.StringType.get()) ); Table table = new HadoopTables().create(schema, Parquet.DEFAULT_SCHEMA); try (DataFileWriter<GenericData.Record> writer = Parquet.writeDataFileWriter(table) .createWriterFunc(Parquet.writeDataFileWriterFunc(table)) .outputFile(table.location() + "/data.parquet") .build()) { GenericData.Record record1 = new GenericData.Record(schema); record1.put("id", 1); record1.put("name", "John"); writer.write(record1); GenericData.Record record2 = new GenericData.Record(schema); record2.put("id", 2); record2.put("name", "Jane"); writer.write(record2); } Table table = new HadoopTables().load("path/to/table"); Expression filter = Expressions.equal("name", "John"); Iterable<GenericData.Record> records = table .scan() .filter(filter) .as(Parquet.readRecords(table)) .build(); for (GenericData.Record record : records) { System.out.println(record); }


上一篇:
下一篇:
切换中文