【问题标题】:Spark write to file with append to s3 - Cost issueSpark写入文件并附加到s3 - 成本问题
【发布时间】:2018-02-09 10:22:26
【问题描述】:

所以我有一个 Apache Spark 流,它每 20 分钟按天和小时写入 S3 parquet 文件分区。
似乎每个批次在写入之前,都会对这个表(/根文件夹)名称的所有文件夹执行“ls”和“head”。

由于我们有多天 X 24 小时 X 不同的表,因此总体上会产生相对较高的 S3 成本。

请注意,我们的架构正在动态更改。

所以我的问题是:

  1. 递归读取所有 parquet 头部是否正确?

  2. 为什么流不缓存这些信息/是否可以缓存?

  3. 您能建议最佳做法吗?

//编写代码:

  withPartition.write()
                .format(format)
                .mode(SaveMode.Append)
                .partitionBy("day","hour")
                .save(path);

看来这个问题与:

https://issues.apache.org/jira/browse/SPARK-20049

Spark partitionBy much slower than without it

【问题讨论】:

    标签: apache-spark amazon-s3 spark-streaming


    【解决方案1】:

    我发现 spark partition by 是这个问题的原因:

    Spark partitionBy much slower than without it

    所以我按如下方式实现它,它解决了问题,而且提高了性能:

     withPartition = withPartition.persist(StorageLevel.MEMORY_AND_DISK());
        Dataset<DayAndHour> daysAndHours = withPartition.map(mapToDayHour(), Encoders.bean(DayAndHour.class)).distinct();
    
        DayAndHour[] collect = (DayAndHour[])daysAndHours.collect();
        Arrays.sort(collect);
        logger.info("found " + collect.length +" different days and hours: "
                + Arrays.stream(collect).map(DayAndHour::toString).collect(Collectors.joining(","))  );
        long time = System.currentTimeMillis();
        for(DayAndHour dayAndHour : collect){
            int day = dayAndHour.getDay();
            int hour = dayAndHour.getHour();
            logger.info("Start filter on " + dayAndHour);
            Dataset<Row> filtered = withPartition.filter(filterDayAndHour(day, hour))
                    .drop("day", hour");
    
                String newPath = path + "/"
                        + "day" +"=" +day +"/"
                        + "hour" +"=" + hour;
    
                long specificPathCount = filtered.count();
                long timeStart = System.currentTimeMillis();
                logger.info("writing " + specificPathCount+  " event to " + newPath  );
    
                filtered.write()
                        .format(format)
                        .mode(SaveMode.Append)
                        .save(newPath);
    
                logger.info("Finish writing partition of " + dayAndHour+  " to "+ newPath+ ". Wrote [" + specificPathCount  +"] events  in " + TimeUtils.tookMinuteSecondsAndMillis(timeStart, System.currentTimeMillis()));
     }
        logger.info("Finish writing " + path+  ". Wrote [" + cnt  +"] events  in " + MinuteTimeUtils.tookMinuteSecondsAndMillis(time, System.currentTimeMillis()));
        withPartition.unpersist();
    
    private static  MapFunction<Row, DayAndHour> mapToDayHour() {
        return new MapFunction<Row, DayAndHour>() {
            @Override
            public DayAndHour call(Row value) throws Exception {
                int day = value.getAs("day");
                int hour = value.getAs(hour");
                DayAndHour dayAndHour = new DayAndHour();
                dayAndHour.setDay(day);
                dayAndHour.setHour(hour);
                return dayAndHour;
            }
        };
    }
    
    private static  FilterFunction<Row> filterDayAndHour(int day, int hour) {
        return new FilterFunction<Row>() {
            @Override
            public boolean call(Row value) throws Exception {
                int cDay = value.getAs("day");
                int cHour = value.getAs(hour");
    
                return day == cDay && hour == cHour;
            }
        };
    }
    

    //还有一个POJO

    public class DayAndHour implements Serializable , Comparable<DayAndHour>{
    
        private int day;
        private int hour;
    
        public int getDay() {
            return day;
        }
    
        public void setDay(int day) {
            this.day = day;
        }
    
        public int getHour() {
            return hour;
        }
    
        public void setHour(int hour) {
            this.hour = hour;
        }
    
        @Override
        public boolean equals(Object o) {
            if (this == o) return true;
            if (o == null || getClass() != o.getClass()) return false;
    
            DayAndHour that = (DayAndHour) o;
    
            if (day != that.day) return false;
            return hour == that.hour;
        }
    
        @Override
        public int hashCode() {
            int result = day;
            result = 31 * result + hour;
            return result;
        }
    
        @Override
        public String toString() {
            return "(" +
                    "day=" + day +
                    ", hour=" + hour +
                    ')';
        }
    
        @Override
        public int compareTo(DayAndHour dayAndHour) {
            return Integer.compare((day * 100) + hour, (dayAndHour.day * 100) + dayAndHour.hour);
        }
    }
    

    【讨论】:

      猜你喜欢
      • 2017-06-26
      • 2019-12-23
      • 2017-10-03
      • 2020-11-21
      • 1970-01-01
      • 1970-01-01
      • 2011-04-03
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多