【问题标题】:Slow chunk response in Play 2.2Play 2.2 中的块响应缓慢
【发布时间】:2013-10-17 12:26:58
【问题描述】:

在我基于 play-framework 的 Web 应用程序中,用户可以下载 csv 或 json 格式的不同数据库表的所有行。表相对较大(100k+ 行),我正在尝试使用 Play 2.2 中的分块流回结果。

但是问题是尽管 println 语句显示行被写入 Chunks.Out 对象,但它们并没有出现在客户端!如果我限制发回的行,它会起作用,但是如果我尝试发回所有行并导致超时或服务器内存不足,它也会在开始时有很大的延迟。

我使用 Ebean ORM 并且表被索引并且从 psql 查询不会花费太多时间。有谁知道可能是什么问题?

非常感谢您的帮助!

这是其中一个控制器的代码:

@SecureSocial.UserAwareAction
public static Result showEpex() {

    User user = getUser();
    if(user == null || user.getRole() == null)
        return ok(views.html.profile.render(user, Application.NOT_CONFIRMED_MSG));

    DynamicForm form = DynamicForm.form().bindFromRequest();
    final UserRequest req = UserRequest.getRequest(form);

    if(req.getFormat().equalsIgnoreCase("html")) {
        Page<EpexEntry> page = EpexEntry.page(req.getStart(), req.getFinish(), req.getPage());
        return ok(views.html.epex.render(page, req));
    }

    // otherwise chunk result and send back
    final ResultStreamer<EpexEntry> streamer = new ResultStreamer<EpexEntry>();
    Chunks<String> chunks = new StringChunks() {
            @Override
            public void onReady(play.mvc.Results.Chunks.Out<String> out) {

                Page<EpexEntry> page = EpexEntry.page(req.getStart(), req.getFinish(), 0);
                ResultStreamer<EpexEntry> streamer = new ResultStreamer<EpexEntry>();
                streamer.stream(out, page, req);
            }
    };
    return ok(chunks).as("text/plain");
}

还有主播:

public class ResultStreamer<T extends Entry> {

private static ALogger logger = Logger.of(ResultStreamer.class);

public void stream(Out<String> out, Page<T> page, UserRequest req) {

    if(req.getFormat().equalsIgnoreCase("json")) {
        JsonContext context = Ebean.createJsonContext();
        out.write("[\n");
        for(T e: page.getList())
            out.write(context.toJsonString(e) + ", ");
        while(page.hasNext()) {
            page = page.next();
            for(T e: page.getList())
                out.write(context.toJsonString(e) + ", ");
        }
        out.write("]\n");
        out.close();
    } else if(req.getFormat().equalsIgnoreCase("csv")) {
        for(T e: page.getList())
            out.write(e.toCsv(CSV_SEPARATOR) + "\n");
        while(page.hasNext()) {
            page = page.next();
            for(T e: page.getList())
                out.write(e.toCsv(CSV_SEPARATOR) + "\n");
        }
        out.close();
    }else {
        out.write("Invalid format! Only CSV, JSON and HTML can be generated!");
        out.close();
    }
}


public static final String CSV_SEPARATOR = ";";
} 

还有模特:

@Entity
@Table(name="epex")
public class EpexEntry extends Model implements Entry {

    @Id
    @Column(columnDefinition = "pg-uuid")
    private UUID id;
    private DateTime start;
    private DateTime finish;
    private String contract;
    private String market;
    private Double low;
    private Double high;
    private Double last;
    @Column(name="weight_avg")
    private Double weightAverage;
    private Double index;
    private Double buyVol;
    private Double sellVol;

    private static final String START_COL = "start";
    private static final String FINISH_COL = "finish";
    private static final String CONTRACT_COL = "contract";
    private static final String MARKET_COL = "market";
    private static final String ORDER_BY = MARKET_COL + "," + CONTRACT_COL + "," + START_COL;

    public static final int PAGE_SIZE = 100;

    public static final String HOURLY_CONTRACT = "hourly";
    public static final String MIN15_CONTRACT = "15min";

    public static final String FRANCE_MARKET = "france";
    public static final String GER_AUS_MARKET = "germany/austria";
    public static final String SWISS_MARKET = "switzerland";

    public static Finder<UUID, EpexEntry> find = 
            new Finder(UUID.class, EpexEntry.class);

    public EpexEntry() {
    }

    public EpexEntry(UUID id, DateTime start, DateTime finish, String contract,
            String market, Double low, Double high, Double last,
            Double weightAverage, Double index, Double buyVol, Double sellVol) {
        this.id = id;
        this.start = start;
        this.finish = finish;
        this.contract = contract;
        this.market = market;
        this.low = low;
        this.high = high;
        this.last = last;
        this.weightAverage = weightAverage;
        this.index = index;
        this.buyVol = buyVol;
        this.sellVol = sellVol;
    }

    public static Page<EpexEntry> page(DateTime from, DateTime to, int page) {

        if(from == null && to == null)
            return find.order(ORDER_BY).findPagingList(PAGE_SIZE).getPage(page);
        ExpressionList<EpexEntry> exp = find.where();
        if(from != null)
            exp = exp.ge(START_COL, from);
        if(to != null)
            exp = exp.le(FINISH_COL, to.plusHours(24));
        return exp.order(ORDER_BY).findPagingList(PAGE_SIZE).getPage(page);
    }

    @Override
    public String toCsv(String s) {
        return id + s + start + s + finish + s + contract + 
                s + market + s + low + s + high + s + 
                last + s + weightAverage + s + 
                index + s + buyVol + s + sellVol;   
    }

【问题讨论】:

    标签: java playframework-2.0 ebean chunking


    【解决方案1】:

    1. 大多数浏览器在显示任何结果之前都会等待 1-5 kb 的数据。您可以使用命令curl http://localhost:9000 检查 Play Framework 是否真的发送数据。

    2.你创建了两次流媒体,先删除final ResultStreamer&lt;EpexEntry&gt; streamer = new ResultStreamer&lt;EpexEntry&gt;();

    3. - 您使用 Page 类来检索大型数据集 - 这是不正确的。实际上你做了一个大的初始请求,然后每次迭代一个请求。这很慢。使用简单的 findIterate()。

    将此添加到 EpexEntry(根据需要随意更改)

    public static QueryIterator<EpexEntry> all() {
        return find.order(ORDER_BY).findIterate();
    }
    

    您的新流方法实现:

    public void stream(Out<String> out, QueryIterator<T> iterator, UserRequest req) {
    
        if(req.getFormat().equalsIgnoreCase("json")) {
            JsonContext context = Ebean.createJsonContext();
            out.write("[\n");
            while (iterator.hasNext()) {
                out.write(context.toJsonString(iterator.next()) + ", ");
            }
            iterator.close(); // its important to close iterator
            out.write("]\n");
            out.close();
        } else // csv implementation here
    

    还有你的 onReady 方法:

                QueryIterator<EpexEntry> iterator = EpexEntry.all();
                ResultStreamer<EpexEntry> streamer = new ResultStreamer<EpexEntry>();
                streamer.stream(new BuffOut(out, 10000), iterator, req); // notice buffering here
    

    4. 另一个问题是 - 你打电话给Out&lt;String&gt;.write() 太频繁了。调用write() 意味着服务器需要立即向客户端发送新的数据块。 Out&lt;String&gt;.write() 的每次调用都会产生大量开销。

    出现开销是因为服务器需要将响应包装到分块结果中 - 每条消息 Chunked response Format 6-7 个字节。由于您发送小消息,因此开销很大。 此外,服务器需要将您的回复包装在 TCP 数据包中,该数据包的大小远非最佳。 而且,服务器需要执行一些内部动作来发送一个块,这也需要一些资源。因此,下载带宽将远非最佳。

    这是一个简单的测试:将 10000 行文本 TEST0 分块发送到 TEST9999。这在我的计算机上平均需要 3 秒。但是使用缓冲需要 65 毫秒。此外,下载大小为 136 kb 和 87.5 kb。

    缓冲示例:

    控制器

    public class Application extends Controller {
        public static Result showEpex() {
            Chunks<String> chunks = new StringChunks() {
                @Override
                public void onReady(play.mvc.Results.Chunks.Out<String> out) {
                    new ResultStreamer().stream(out);
                }
            };
            return ok(chunks).as("text/plain");
        }
    }
    

    新的 BuffOut 类。这很愚蠢,我知道

    public class BuffOut {
        private StringBuilder sb;
        private Out<String> dst;
    
        public BuffOut(Out<String> dst, int bufSize) {
            this.dst = dst;
            this.sb = new StringBuilder(bufSize);
        }
    
        public void write(String data) {
            if ((sb.length() + data.length()) > sb.capacity()) {
                dst.write(sb.toString());
                sb.setLength(0);
            }
            sb.append(data);
        }
    
        public void close() {
            if (sb.length() > 0)
                dst.write(sb.toString());
            dst.close();
        }
    }
    

    此实现的下载时间为 3 秒,大小为 136 kb

    public class ResultStreamer {
        public void stream(Out<String> out) {
        for (int i = 0; i < 10000; i++) {
                out.write("TEST" + i + "\n");
            }
            out.close();
        }
    }
    

    此实现的下载时间为 65 毫秒,大小为 87.5 kb

    public class ResultStreamer {
        public void stream(Out<String> out) {
            BuffOut out2 = new BuffOut(out, 1000);
            for (int i = 0; i < 10000; i++) {
                out2.write("TEST" + i + "\n");
            }
            out2.close();
        }
    }
    

    【讨论】:

    • 感谢您的回答维克托。缓冲将提高速度,但是从我写出到它出现在浏览器中之间的延迟仍然很大。添加简单的 println 语句表明所有行都将被写入 out 并且当没有更多行并且 out 关闭时,它们开始在浏览器中加载!如果行数太大,就会出现这样的超时错误:
    • [ERROR] [10/22/2013 13:57:16.285] [application-akka.actor.default-dispatcher-5] [ActorSystem(application)] 无法运行终止回调,由于[期货在 [5000 毫秒] 后超时] java.util.concurrent.TimeoutException:期货在 scala.concurrent.impl.Promise$DefaultPromise.ready(Promise.scala:96) 在 scala.concurrent 的 [5000 毫秒] 之后超时。 impl.Promise$DefaultPromise.result(Promise.scala:100) at scala.concurrent.Await$$anonfun$result$1.apply(package.scala:107) at akka.dispatch.MonitorableThreadFactory$AkkaForkJoinWorkerThread$$anon$
    • 您能在您的代码中插入几个System.out.println(System.currentTimeMillis()) 并在此处显示输出吗?请将它们放在static Result showEpex() 之后、// otherwise chunk result and send back 行之后、public void stream(Out&lt;String&gt; out, Page&lt;T&gt; page, UserRequest req) 的最后一行之前以及return ok(chunks).as("text/plain"); 之前?由于某种原因,您的块执行没有完成或花费太多时间,因此执行被播放框架终止。另外,您是否尝试运行我的代码?您能否确认您是否有同样的问题?
    • 这里是输出:1382542832461:showEpex() 1382542837811:showEpex() 1382542837875:开始分块 1382542837988:返回块 1382542882427:退出流() 13825423第二个退出流不是第一个!如果我只使用您的代码而不是对结果进行分页,但使用更大的循环(如 10000000),那么当我流式传输大(16GB)数据库(超时)时会发生同样的事情!似乎不是生成和写入,而是首先生成所有内容,然后生成已经收集和生成的数据块!
    • 更新了答案。应该可以帮助您解决问题
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2014-04-29
    • 1970-01-01
    • 2021-08-17
    • 2023-03-12
    • 2013-02-16
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多