Return-Path: X-Original-To: archive-asf-public-internal@cust-asf2.ponee.io Delivered-To: archive-asf-public-internal@cust-asf2.ponee.io Received: from cust-asf.ponee.io (cust-asf.ponee.io [163.172.22.183]) by cust-asf2.ponee.io (Postfix) with ESMTP id DB28F200B84 for ; Tue, 20 Sep 2016 10:35:32 +0200 (CEST) Received: by cust-asf.ponee.io (Postfix) id D9B05160AC9; Tue, 20 Sep 2016 08:35:32 +0000 (UTC) Delivered-To: archive-asf-public@cust-asf.ponee.io Received: from mail.apache.org (hermes.apache.org [140.211.11.3]) by cust-asf.ponee.io (Postfix) with SMTP id BC846160AA9 for ; Tue, 20 Sep 2016 10:35:31 +0200 (CEST) Received: (qmail 81109 invoked by uid 500); 20 Sep 2016 08:35:30 -0000 Mailing-List: contact user-help@flink.apache.org; run by ezmlm Precedence: bulk List-Help: List-Unsubscribe: List-Post: List-Id: Reply-To: user@flink.apache.org Delivered-To: mailing list user@flink.apache.org Received: (qmail 81099 invoked by uid 99); 20 Sep 2016 08:35:30 -0000 Received: from pnap-us-west-generic-nat.apache.org (HELO spamd3-us-west.apache.org) (209.188.14.142) by apache.org (qpsmtpd/0.29) with ESMTP; Tue, 20 Sep 2016 08:35:30 +0000 Received: from localhost (localhost [127.0.0.1]) by spamd3-us-west.apache.org (ASF Mail Server at spamd3-us-west.apache.org) with ESMTP id 4A3D01805E6 for ; Tue, 20 Sep 2016 08:35:30 +0000 (UTC) X-Virus-Scanned: Debian amavisd-new at spamd3-us-west.apache.org X-Spam-Flag: NO X-Spam-Score: 1.179 X-Spam-Level: * X-Spam-Status: No, score=1.179 tagged_above=-999 required=6.31 tests=[DKIM_SIGNED=0.1, DKIM_VALID=-0.1, DKIM_VALID_AU=-0.1, HTML_MESSAGE=2, RCVD_IN_DNSWL_LOW=-0.7, RCVD_IN_MSPIKE_H3=-0.01, RCVD_IN_MSPIKE_WL=-0.01, SPF_PASS=-0.001] autolearn=disabled Authentication-Results: spamd3-us-west.apache.org (amavisd-new); dkim=pass (2048-bit key) header.d=gmail.com Received: from mx1-lw-us.apache.org ([10.40.0.8]) by localhost (spamd3-us-west.apache.org [10.40.0.10]) (amavisd-new, port 10024) with ESMTP id 2TsyO-Y6lBsD for ; Tue, 20 Sep 2016 08:35:28 +0000 (UTC) Received: from mail-vk0-f54.google.com (mail-vk0-f54.google.com [209.85.213.54]) by mx1-lw-us.apache.org (ASF Mail Server at mx1-lw-us.apache.org) with ESMTPS id 2DECB5F47C for ; Tue, 20 Sep 2016 08:35:28 +0000 (UTC) Received: by mail-vk0-f54.google.com with SMTP id o139so15285699vka.1 for ; Tue, 20 Sep 2016 01:35:28 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20120113; h=mime-version:in-reply-to:references:from:date:message-id:subject:to; bh=Mop3BEIQYQSy36dwjYVK5Icte8j48o8uAE8hqLrd9u4=; b=zfmn7zgU+gDVuSlWIUJt3RJmjFYuKMNiOgVFUBnGdykhMT/zN4jI9PiEdyyOufrrb/ wyL6a3T5gnXbN8AmJVBTF+toOmZkyQnr/CIs5MgUOAJBNCuqKq8HSZnRFBwKH+vcD81O GOZo3E2s/Z3xYut09ZOSoT9vuOmmtxvi+8/KRYXFr+r6IjUCwSnD7Qwal1xtHILhKSEG dxwj4vO42w+UPm0twfFf+vlzTYa8MWKcCWWAiJr7hY0ABndxxxwYJL6zLt4qVhyn9v6K B5azU/Mco5H6ulPMRYLIwaIDxFtixRISRIFQRgWAPJTKo/lCAPpY8Qo3PtOZWPUz5G6V 3Axw== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20130820; h=x-gm-message-state:mime-version:in-reply-to:references:from:date :message-id:subject:to; bh=Mop3BEIQYQSy36dwjYVK5Icte8j48o8uAE8hqLrd9u4=; b=HDlqF0IQuD2OTDFzBRUSMdgyLhYAwAMubdiTJq7K+XQRnbFBhAYdwciUa3NCrWs724 pB7EcWNqWvCzUdDxoscydlxPcNUVso7WSPlennfmaPp52aYPeIcr3GJ7rDApxk9aO/lR TW+PL2GpqUpEicMeInQrqGBM8GRPaTfo4wmLrCmWH6OdUm4YvSNdo24dwvZ/Bi0f7xLZ 4s6tGS/FyOG4+u/tLomlhtEvcFdQaLZk27HxJEGgNPWhcmyhc+q+TcqikSefGewum8pJ EmPpEgZVG4OAX3KjVIwXPgTFjd2ebaGqof5T4yui+gEsD+TKk1+nyMaf6bTPBDU8gZ21 sExQ== X-Gm-Message-State: AE9vXwOGn15j8Z/3IMCUncGLrX6+3ztZd7VHfjdu8HLy69JzdESxFImEkiwZBjI+rItUFRLLBB11dW7sofVNRw== X-Received: by 10.31.88.135 with SMTP id m129mr2581005vkb.159.1474360527435; Tue, 20 Sep 2016 01:35:27 -0700 (PDT) MIME-Version: 1.0 Received: by 10.103.153.193 with HTTP; Tue, 20 Sep 2016 01:35:26 -0700 (PDT) In-Reply-To: References: From: Yukun Guo Date: Tue, 20 Sep 2016 16:35:26 +0800 Message-ID: Subject: Re: Serialization problem for Guava's TreeMultimap To: user@flink.apache.org Content-Type: multipart/alternative; boundary=001a114e1aa2823f14053cec4f17 archived-at: Tue, 20 Sep 2016 08:35:33 -0000 --001a114e1aa2823f14053cec4f17 Content-Type: text/plain; charset=UTF-8 Some detail: if running the FoldFunction on a KeyedStream, everything works fine. So it must relate to the way WindowedStream handles type extraction. In case any Flink experts would like to reproduce it, I have created a repo on Github: github.com/gyk/flink-multimap On 20 September 2016 at 10:33, Yukun Guo wrote: > Hi, > > The same error occurs after changing the code, unfortunately. > > BTW, registerTypeWithKryoSerializer requires the 2nd argument to be a `T > serializer` where T extends Serializer & Serializable, so I pass a > custom GenericJavaSerializer, but I guess this doesn't matter much. > > > On 19 September 2016 at 18:02, Stephan Ewen wrote: > >> Hi! >> >> Can you use "env.getConfig().registerTypeWithKryoSerializer( >> TreeMultimap.class, JavaSerializer.class)" ? >> >> Best, >> Stephan >> >> >> On Sun, Sep 18, 2016 at 12:53 PM, Yukun Guo wrote: >> >>> Here is the code snippet: >>> >>> windowedStream.fold(TreeMultimap.create(), new FoldFunction, TreeMultimap>() { >>> @Override >>> public TreeMultimap fold(TreeMultimap topKSoFar, >>> Tuple2 itemCount) throws Exception { >>> String item = itemCount.f0; >>> Long count = itemCount.f1; >>> topKSoFar.put(count, item); >>> if (topKSoFar.keySet().size() > topK) { >>> topKSoFar.removeAll(topKSoFar.keySet().first()); >>> } >>> return topKSoFar; >>> } >>> }); >>> >>> >>> The problem is when fold function getting called, the initial value has >>> lost therefore it encounters a NullPointerException. This is due to failed >>> type extraction and serialization, as shown in the log message: >>> "INFO TypeExtractor:1685 - No fields detected for class >>> com.google.common.collect.TreeMultimap. Cannot be used as a PojoType. >>> Will be handled as GenericType." >>> >>> I have tried the following two ways to fix it but neither of them worked: >>> >>> 1. Writing a class TreeMultimapSerializer which extends Kryo's >>> Serializer, and calling `env.addDefaultKryoSerializer(TreeMultimap.class, >>> new TreeMultimapSerializer()`. The write/read methods are almost >>> line-by-line translations from TreeMultimap's own implementation. >>> >>> 2. TreeMultimap has implemented Serializable interface so Kryo can fall >>> back to use the standard Java serialization. Since Kryo's JavaSerializer >>> itself is not serializable, I wrote an adapter to make it fit the >>> "addDefaultKryoSerializer" API. >>> >>> Could you please give me some working examples for custom Kryo >>> serialization in Flink? >>> >>> >>> Best regards, >>> Yukun >>> >>> >> > --001a114e1aa2823f14053cec4f17 Content-Type: text/html; charset=UTF-8 Content-Transfer-Encoding: quoted-printable
Some detail: if running the FoldFunction on a KeyedSt= ream, everything works fine. So it must relate to the way WindowedStream ha= ndles type extraction.

In case any Flink experts would like to= reproduce it, I have created a repo on Github: github.com/gyk/flink-multimap

On 20 September 2016 at 10= :33, Yukun Guo <gyk.net@gmail.com> wrote:
Hi,

The same error occurs after cha= nging the code, unfortunately.

BTW, registerTypeWithKryoSerializer r= equires the 2nd argument to be a `T serializer` where T extends Serializer&= lt;?> & Serializable, so I pass a custom GenericJavaSerializer<T&= gt;, but I guess this doesn't matter much.


On 19 September 2016 at 18:02, Stephan Ewen <<= a href=3D"mailto:sewen@apache.org" target=3D"_blank">sewen@apache.org&g= t; wrote:
Hi!
Can you use "env.getConfig().registerTypeWithKryo= Serializer(TreeMultimap.class,=C2=A0JavaSerializer.class)" ?=

Best,
Stephan


On Sun= , Sep 18, 2016 at 12:53 PM, Yukun Guo <gyk.net@gmail.com> wr= ote:
Here is the code snippet:

=
windowedStream=
.fold(TreeMultimap.<=
Long, String>create(), new FoldFunction<Tuple2<String, Long>, TreeMultimap<Long<=
/span>, String>>() {
@Override
= public TreeMultimap<Long
, String> fold(= TreeMultimap<Long, Strin= g> topKSoFar,
= = Tuple2<String, Long> itemCount) throws Exception {
String item =3D itemCount.f0;
Long c= ount =3D itemCount.f1; topKSoFar.put(coun= t, item);
if= (topKSoFar.keySet()= .size() > topK) {
topKSoFar= .removeAll(topKSoFar.k= eySet().first());
= }
return topKSoFar;
}
})= ;

The problem is when fold function= getting called, the initial value has lost therefore it encounters a NullP= ointerException. This is due to failed type extraction and serialization, a= s shown in the log message:
"INFO=C2=A0 TypeExtractor:1685 - No fields detected = for class com.google.common.collect.TreeMultimap. Cannot be used as a = PojoType. Will be handled as GenericType."

I have tried the fol= lowing two ways to fix it but neither of them worked:

1. Writing a c= lass TreeMultimapSerializer which extends Kryo's Serializer<T>, a= nd calling `env.addDefaultKryoSerializer(TreeMultimap.class, new TreeM= ultimapSerializer()`. The write/read methods are almost line-by-line transl= ations from TreeMultimap's own implementation.

2. TreeMultimap = has implemented Serializable interface so Kryo can fall back to use the sta= ndard Java serialization. Since Kryo's JavaSerializer itself is not ser= ializable, I wrote an adapter to make it fit the "addDefaultKryoSerial= izer" API.

Could you please give me some working examples for = custom Kryo serialization in Flink?


Best regards,
Yukun




--001a114e1aa2823f14053cec4f17--