From c7ad532fb2db7f431f061f685f82863f1d196eab Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Mon, 7 Mar 2016 22:47:33 -0500 Subject: [PATCH] GH-13: Add checkstype and fix as much as possible Fixes GH-13 (https://github.com/spring-projects/spring-kafka/issues/13) Also upgrade to Gradle 2.11 to support the latest `checkstyle` --- build.gradle | 12 +- gradle/wrapper/gradle-wrapper.jar | Bin 52271 -> 53638 bytes gradle/wrapper/gradle-wrapper.properties | 4 +- gradlew.bat | 2 +- .../kafka/core/BrokerAddress.java | 3 +- .../kafka/rule/KafkaEmbedded.java | 47 ++--- .../springframework/kafka/rule/KafkaRule.java | 3 +- ...kaListenerAnnotationBeanPostProcessor.java | 16 +- .../kafka/annotation/KafkaListeners.java | 1 + .../kafka/annotation/TopicPartition.java | 6 +- ...AbstractKafkaListenerContainerFactory.java | 2 +- .../kafka/core/ConsumerFactory.java | 1 + .../core/DefaultKafkaConsumerFactory.java | 1 + .../core/DefaultKafkaProducerFactory.java | 1 + .../kafka/core/KafkaException.java | 1 + .../kafka/core/KafkaOperations.java | 2 +- .../kafka/core/KafkaTemplate.java | 12 +- .../kafka/core/ProducerFactory.java | 1 + .../AbstractKafkaListenerEndpoint.java | 8 +- .../AbstractMessageListenerContainer.java | 7 +- .../AcknowledgingMessageListener.java | 2 +- .../kafka/listener/Acknowledgment.java | 2 +- .../ConcurrentMessageListenerContainer.java | 12 +- .../kafka/listener/ErrorHandler.java | 3 +- .../KafkaListenerEndpointRegistrar.java | 2 +- .../KafkaListenerEndpointRegistry.java | 4 +- .../KafkaMessageListenerContainer.java | 80 ++++----- .../ListenerExecutionFailedException.java | 1 + .../kafka/listener/LoggingErrorHandler.java | 3 +- .../kafka/listener/MessageListener.java | 3 +- .../listener/MethodKafkaListenerEndpoint.java | 2 +- .../MultiMethodKafkaListenerEndpoint.java | 1 + .../AbstractAdaptableMessageListener.java | 2 +- .../adapter/DelegatingInvocableHandler.java | 7 +- .../listener/adapter/HandlerAdapter.java | 1 + .../support/converter/MessageConverter.java | 1 + .../converter/MessagingMessageConverter.java | 2 +- .../EnableKafkaIntegrationTests.java | 13 +- .../kafka/core/BrokerAddress.java | 3 +- src/checkstyle/checkstyle-header.txt | 17 ++ src/checkstyle/checkstyle-suppressions.xml | 8 + src/checkstyle/checkstyle.xml | 169 ++++++++++++++++++ src/reference/asciidoc/quick-tour.adoc | 24 +-- 43 files changed, 343 insertions(+), 149 deletions(-) create mode 100644 src/checkstyle/checkstyle-header.txt create mode 100644 src/checkstyle/checkstyle-suppressions.xml create mode 100644 src/checkstyle/checkstyle.xml diff --git a/build.gradle b/build.gradle index bd0f31a5..8c70ad80 100644 --- a/build.gradle +++ b/build.gradle @@ -47,6 +47,7 @@ subprojects { subproject -> apply plugin: 'eclipse' apply plugin: 'idea' apply plugin: 'jacoco' + apply plugin: 'checkstyle' if (project.hasProperty('platformVersion')) { apply plugin: 'spring-io' @@ -100,6 +101,11 @@ subprojects { subproject -> } } + checkstyle { + configFile = new File(rootDir, "src/checkstyle/checkstyle.xml") + toolVersion = "6.16.1" + } + jacocoTestReport { reports { xml.enabled false @@ -325,9 +331,3 @@ task dist(dependsOn: assemble) { group = 'Distribution' description = 'Builds -dist, -docs distribution archives.' } - -task wrapper(type: Wrapper) { - description = 'Generates gradlew[.bat] scripts' - gradleVersion = '2.5' - distributionUrl = "http://services.gradle.org/distributions/gradle-${gradleVersion}-all.zip" -} diff --git a/gradle/wrapper/gradle-wrapper.jar b/gradle/wrapper/gradle-wrapper.jar index 30d399d8d2bf522ff5de94bf434a7cc43a9a74b5..5ccda13e9cb94678ba179b32452cf3d60dc36353 100644 GIT binary patch delta 12684 zcmZvC1yr2DvLVQ_bc;O_43F2UU)NS2W7yPJ1+&Y9Ca zUw>6y-Cci8O;3GK1P{ysN0gHUhkyeCfq?c9fwZ6-Z=GsHFnh2+kUzxkAIKUq=MPkk0fP#K z@>_lcK}RGS33Iv1!7!835ltIJ;0Ig)v1d#iGE$_|w@%j2>XOo_gq-Jigz#=Js zP=?14^As$%jVIfQSkY#^P&6m~a1lV%fn=Q8s+n23+{4aN&2p9Te4_REI6kfmvU^S) zB$+s<#}v!-N1U~Q;jWd%4?HkyP3fv5Raec@ITCwH_jub3H8Oi`hGn_Z zk0Q{uwr(^^&FQWv-DV-;rp{xgs^!k41a0l)G%3w$y81z#9If-~3Cj#_`8?ZFdgFR! z*m$~|r=?q;%nw!Vy`0rkcS=(dz>4QG+}!(?7P+!PoAF0`#Y*?epz8QM`N?V8g6;R^ zMaLz{LB+pri$z46V@!n>QHttKH_5!$6(W>~WcWqL&Ul65%o+-V1*cEkJ?7q}oOvtN zZm7TDB4EHz>gZa(|IT{YtU3B~_p~IO-3h^OfCWC)o8*P=f)zb#UT!OY1P~qwiXgnk zag(7NSI(uoC2};e^o|MiB?v?z_~ILS@BK*@BaAM38@VUvHMm}C#xYZxGxhjDyJ) zfavbgVnr-+YADcj%T5bz3qWxgabfNC6C%kWRx2@$cy!dZ%XWr>@?7lQ`j9OFJ36r) zj*$W^v%ogSYcE#>w-pKI7t_<0rbJXOJaZ02V;PkA5V1pU(iyk{GBHphsuaE|6QVH( zYH)}+Qv50xNle~HibgeSs2X?$GF9LGFRz^{ndEt| zPeWivTr(@r@BV!Ns^^dyA^*5FEN0|iw@$@YhDQAD)*Neft9WQ2AjdSoUqsY^cb@1w zwft1gi;MK#bjV#p7DA+MLF~B52;xx-$WiA^#D_*MF=5AKr(-cfWg4Lpvu$#klPaZ= z1`Nvb;uF*hU%cvWvubN=XH_h0j!s%z?{3$YA8GFZCeG|BM#x)??awB=9;ti3QqMef zb)FxmMZcMm0FX2J#uECd!b!scy2-v^A3@Fg1`BXAnF+t2Lm*wGj6+g)pcz}nd<24P z4JhU&;ck$oV`g$Us3jl3L4$vcGaYzJf*ZfQItCfnm)@KPcf=h>ILG0Ub(hHHRNQ<3 zP43WjNG`V18fArJJiq!JfHfG_3z|GGnF}b&NEmbgGzNRbhouo@ z8kFB$!MT;wZ#cLr`;6~gfw`6_x;pv9`(gDR=9x$3%Tt%4!0kP{m0f?N@%-AAap_Q< zF)uDE=`*%{Gi3>L9?&AcIR{dhD8?E} zrpo4WSHk^+!b5{pGHt*O0HoHgxho6AYR{QJ=-KkVQSiqJ{l- z`cDKEPtno4M zHEL{HB>EzRTO6{n+cec=6**1!<56c@O(Nb-z!H7$tNOP<7+b)0jgyteHrJP_YsT@}gBUnK020C23K zh4$PqO9M0em zF5dj*s@Vp4GoDSbD8N9{}COi!GI~iXrv|HO>d2 z@$W69g_9^ng>C{<1p>Ycx?BxdQ`XX8f}aH==5#Sh%6z$Y%L)F3EDT}Jgr<%YB$c^T zWy#XKgO?9!8D^~>uXN9HjtV7M1XG$u_G@KoQ)y_$cW~PHnz{0uE*k!_K5#sg^s@=) zG*E;(j;X9I*$Lo7X}hTBj4Ljw;tT^7QpJ;OuuAjJ5rH>EK2X{7q$mPWz7N+f% z&qtu%F&~V{iCiz4PojACmUbd>(Ok``@*I;)-}!oH>5Cx`q8I~q*1Fk^2uTsg{49b9 zb1YLc0g`*S3D}{Oy^0=>!nEa*y9^#7sMl%|^@PZXHYXr;nO41dIGUFr>_%I}%v8a% zy2i1lhH4hdu2_hGu_MrZK;i&pc^Y!il@^8B(lwIBy*WI}*#86t(DM~tX?kTds-##| zmVB4Qq%N{1h25ZsDatgP#*CHQuqjRgJ4c%-CO=d>p9-O*5-YQ+Cj>7{Xhw<>av-=g zdZrF$*a(10w(nA*@6V3wz91egPKPh6u1>sEI^TmPzgl$+y9T3xyYE;}#=D2@BuM&x z(AjAyj3c%pYch5c&8P=E@akr9Lo*ukYil+QRQxw3_OJ|YiVbq#+6(IaFc-eV?a4wy z!Pp`-6xq^Y_gN`xR>Ot-%-Vt)OFZ;2Dosaqbz}e=X<+WOe`p?)ymo*sO#& zrw3q&45yS8qQ{!B3LVp|&;wPhK0=INJnc&jxpc=8tCmX^7A(btqq`er?jCY_9G>IA z`5fv$F;s~*RQ8p*1~VgdPTWBnXP%9ob0mIPQ z#k$5wjDe`SuBu^HW)E>|_PXUuQvTG{wJ$)TJt~|>^ z+1;N|`9P$NFsxt6b|`zixF*F~?!f(LX;-sy+gBT<2gmer( zJrbkfrQ-Wu0y~VC37Nv>F);gcwNA=~is&y@Lmm%o5YDImskV?P?gE?3%58q_*_Qy- z2bs9L3+;jTr^bF=E+tQ_FHJ*t=iYH=J2qkt8AErWs7QrR#Lw_}ofBK(MTy2D;g(${ zO~;;^R1Ebh6T1oIdls7w(K7~il7a$X4osw`4ZRZMJG}7{@KwrK^VxMrLo5iyek4zZ zdT~t^FpP{Hm53yL=mB^=?(C~Pq<8{8!QVy=t#3U`am!E(CA*-~Z)scu#WOFvpRD3e zNfP&o0irvidKjlwbPOy;1=C;J$~>CYI6* zAAjt(CHVk+#R=mQ=RPIxXkFAIo%r22eas`qPe)4#oVJSR-PDRlZcoN&= zwgp@MI)+$_?@&N{vY!gFVRE%yBUQ~ir(OcyfPGpmnrZW?Y;b>IENZ9h1Vp%K5y15zvN|Z zt3Ju}P2~)Ii~}GLUD-eYxLHKG44*50TDe#A-io*R;*;&81p=}>` zA+1k$)%^63au6LQI$Cdb480R(YIV1B55H3tTa+g)eVGafQH7w!w;RHj0VtR+4siP( zkWud7R~N)L5th~h&HEe=)J}!?Cc{5(!c?Ag3#znc8GQ6{TO1AmtY$P^;#k)zs*JUg zUmYB5OQjtx_K2|EcFDov@SDue2y%Rz(4WU2+l7NP^&HWL8XJ*YpJT-IeR-!nwrfn) z2%oeS(-swR4)cz({`-3iUy-8I3TyEPKSb#KHd_q{8}lTs7%9T^q+^e6I~C892`_PZ zb0FQCCr72DhPE?40EV%$8lTBDb}lkKqt&2@r?T97`dNjCtlxhVgD;jD7(QwJPE<%Y zC@ucPVTAf{S+u!dIfIma@JBwgxehFv`$m;~-$2u{sn%EYRD1!7ZW)=`G5w`)FHN^b zJTX=9-wHNUjEHE#?XcmZl2?JTiIUeG$)TDf6Zrte!1!=JCVv|)4?P2xyBs2N+HJRt5 z1YL-LqE177xha@+I*&#U66b>}l@PM7_Hkq`T6)Go1~Tv%dqyRjs$tWs<9FG!u~}Ky z($UvHDezNH07sc|hM#u}B`F3I7P;{T3Gt@#Xs443jAAfwQaM4%i!|RCvDKUW+KXnX z=o$kZy&xo)LB|#3vEOUBl#>&{Zaw0bL6(uey-}O8Yd>O;ulP6*(QRWKq6LktpzUt2 z&7=QOB#fGw1%&0j$+gD8cr=~~n+ZnZWaGi5unSjxC7_dpJ))_!Xt1`kBD4@~)3_Lq zp~#2K?!y;vB<5y;0!5-Y6nQ?#1Bs@P#PU8>La$`F9oO)?&zu&^vZ*mAR1qVtGMy`U z2b47)=S~FFt{17Qe4ik9{Lszl@Y$1e2#If)!WqAFNJ+Tu=wHUE66C`MSt4c#df+}x zjP*){=L4p2qVx$z^HE#U#MupVSF*(GC(49wQwDl&Q^ev-+rMcMUVEpv`4C3$!a$03 zIw%YCOKSz0C`dj4KY}s>Em{R^_%Y^ohYcmlke;{w+60o9R0yvjujuHCp$ZE#@_9vj zBg3Pr;kn}d>?w@OvT>9U2C)|G%bgmGceI3ca{-J3*`8b}ZBiZr6rUvb*%s#gYIhvD zS7Ap@K8<@hVwsHQB{{OMd0w-Y-Wu%8S3>P6ixOFX)&(JBN0;IM-09@!sG->lGHP}~ ze-Sv_k~g#6SAh32((kIHb>U{L>w>f zat8%5B-N!^gJ$)gk3-E!Ye8|BH@Z$f2EEUJr~wab@*g85*oB#VK+g|1nlPRbpmdtT zFngSEVGy=NbNiQ9sp5)nj?^z1_(3Ge|RcqZKhC9wl)rZgyk zuemVY-N9!qX11mm3ToLVvXj(JtuY>Pay#=5V$c|^z$0Bqwp=Eth`(&Qy28iSs?;yqEOcsivfxl)0#mmcOhG~rp|4$vZ0oty|{K1>pFzC_91c_wR2kG3`K(t zz1qR5Kno#oeu9I;fR^+_uo}<1?2I4)wk61=O`Ty!5MLY=u{WPddU#>WkV^H35)(3B z<2`&gw5PpA8`)*Mn35S{{2$aeXgm~ihfwC66`t`BcZ-i^9|&t1ltmNVc!w`TkrD2q zRg`1f3+J$|Z>B@Ma#f(bWG9DsmYTSU#Jub62iU8veCq|umN!$l3uD`SD_&>;5N&>A z0cuGgjF3xb^?r`D${kIjmEK&o-XkKb@b}l8EjyvM>X1iM-q(J~u8fRFyajvW*Pioi zCwJ!_={mMkH0|K;A@_eF2XmPcxgjAa*l7xH7K3uR5Y)p@xmmTD7L|`GHsF@il?H8Z zv?i|Zk@)jP$hgNxOxsbq@S9P=0_?&FvMZ7YVkA85tpr83HDla zR=^hGc9_>c#5BTV>i8xIPPD>YJR#c$^JIKNwCiw&#mHGjI$D1hx*8M9>$W zFsF3oYT8QJR2}Q?68kyQ4GW)^s;_#Qm-s6u(Q2JD0fL`7E`E{1UjtCpTEp~QAt^mD zXpULQpAyVbX3(?ZeY9asTBY0l7toLiV;p-hxojhka||YP75K6z@-hrQ9pL$)L-#Z<2Zkd| zV1UPd0!6--P-Wi7BL+xrq;q}RCfF$O9jXNHx53J9&YiYthkWiw6v!cPbr+31oTqgT`#$1T9OxWgLtFb7Fs798i+ zr<;hxgR;I(|4DXZ%W%^p+#If;DumLLLNZeC7h+%OO`^9DO##N%SkNDKGuu z5qBHd2TzR8x=SN?dIaLE-ibIj6?;*16a{7%rZf)%gNxFm&lT0xDUSzd-7XI2oMS%@ z$9P=Wh|b2E?N?w+nHp%p=sk-UmE93l`Lw<+8fgE z|BNP(8X+5!7CC&Z$o%LCl7T`Laeo2uKgW+wdo;RL=jX%fS>(X(0A48MFi#^m4t;1$xtz%HEzuS7e|C)H)P3-MY z1bu<_^wB4Srf_`ib;ui~<^fKz>*% zmy(S}ZOSTlDn7o-*=#VZT+cY&r+myTQuOo7uH4=L+xTus%jkaH*ec!_kA;h4jc(_q zcB^1>4w+v!ZvZMUH{_P@bpD636qV1ibZWba>6NI15m_r-qu(_)mXgJ$UTA|@^c`Us z5KpH^(7gyo)TVLhsD?zc9NaYkvc0u-DV01Xxwqt$=%q(2Bi|xm6+@|RBE0fND60(% zu1@OQnDb32!qAask3UG^IwlcT*O5y_1iF3d>>YhjL05I-;pKh2Q48!U?3%agIZ$)b zCCx)J%I|Y>5Q>oG2~Uh)$x4FFS{42}S{q1`oc_t3{%*wz>!TJ_fVO<*Pp=K+^kULzCP0C@p&nn z1(xJ>V~mOd!kz+Z%2N|t0Ks%2D-HByUrpj@1C+xhB^k8w>-@x=u}eyR@6JNm{dg=S z-_E!edZme}r_Ty=-hd7CKEfirXrq)GHm^TybiC)4?ltQTS|kQ8$S|Y0f?+tnhb7>tur{8 zGQY(-grsVg<~Fp64wLUao-~HW)RYt!`V{F#>(76yhI65zUOoygHN?TmpUn^}7_8s0 zWco56rueHC@@D@1B}VDXuQ&;;7Mk+RMxS+(Wvwc&EWwh*9>$_3HR*x1??z&U_pILC zT-vs-9Kb46MfdUw3_wj2TL6}S(A`9$?R2uEB8PR96po78`uLQ&=_vgGE`Gx4e!q<} zPx$(`uJ!6KSn9mA&fCVceHD_ud9IZq2^)E|h#L|ZG}yE+KKu&mD~jK>H8fN&sCaU= zoEpS?GlaDX`=d;xi|<2ipOo2-T2dy(!xPLtvgB>qM1-X%^8n~9597WaQ0HWEl?S%r5RoMMw^f4RrU_*HPb=7$rX-9HjEP`co$Bhj}{QXZNlqX%Df-Sc( z^bWoHkiFW|T-r9%oTmVHj=}U4 zQcNouOWUqcMoUF4iQ$Ndt>4@@fp5X=+34LNpOQXjXkeWRg+7|^M#d-*&O%LUbYV|SggAA zkOMc9o#%Jz8c~w2Y3$JHDo^GSRJeENq3-Q(&CMv}$!Oj5*FDB^3BHj0)rnq^o2XkHA z_Z7tf?E`|}g~hC$#{}&gmvz_9ZWHGD+La4^0`1a+qj5iR%7W7EWrmGr8PTKCrUlp_ z8XHzq+av`5`r`5S)W%pou{cHyrz$QO9z`XY=_~7zQgMT}kIrHalG zb1I7#RGb!6+!p!<&Jvg>K0Y=$w*oHIpcw3TOE_6wlj|=la~vC zXGcSX(=J}MJ+Mj)Nlt5@$fH>ff?5tgd44piP5SawKbSF!T9noWAr7|g9=8V2Jw*|! zJS`s&JCdsXI#3?JiIvcu2KWR;LV?eyifpmfY=t4bKzQ0pQ!pEvJx?6_*up+HnOj^~ zH47qnOAoLc=UpD3WcsDXsZRKzt(qi2e1Kk{rIlvUGihOsBkTJMh?>dT?63#ka5}R9 z@u7a3-b|PU}k znBK79v)#xy)eYP9HbBF`oz?^$EoqzHb(rTP)9mA3zRNRXC$V$LTRbgziSsYV^(g~j zxB?HN8`&2U1|*xBtU+rZz&z~$fOjXY{_SmoFo*;7;cl9rbBg6e`?&gQ-(&g5K`1r} zK+dYj>jy6?_1l1NVuE+G&6L;+w+@Z*A5he~a+X(VyNom?^S7O>yi|SFg*EOdDtqaY ziM#OU-e@YGNU~VDA~=0u=G1k)rybv{TG_O)OA8f`5s^%Qt>s#RX8XAW&^S!G^}E9| zT4(LK)QCTKY6Oe5Ph@OH))BtO@NniOQ!p6U@teQ2j51jllk?;RA%A&d-{PuRWa#~r zE3ALUkv9=^#7luWU=F{5;9M#^C0(9d7OQk9_Jx7+Ktwh(BA&dO-Co6DxMeZU?OhRb zvwVfxA_kk~_3UFrX3H8oVA>tRxQwObR7lwWq18%6eWE_O>Ql3CteUhsNlOypN>}Iw z8|BE!N_)4Tq2XA3Un)a*IRmo?n(Vs>IvV{2k)t_$cfr$}89nx`EarNY$cN)Y$L|GT z_@}qTOHtFS8;2@pIp^yy{QqXg{K*XNcE?nD#9pPUyg@#^-N=d;3>{T1`!wmfg2-%1W>hb!JEZcoyS!V;St9N1(gqiU?Qelkbs9Jm-&pw?IZAkL}sf` zTC;D8IFcn@UX+X~3n{aui);ccFVTev*-R2oqm&~TA=0_5eAfF|)Eg3zX-pJuG+ebs zd85I1x5l!UUVX;-$ag*d%hR5)L*&T|r^!yO$3@6fwdYQV2SD5voCR%&Vq_RhiuDyo z5Z;j4RnGl^x-YL*h=pV;7)GBtu~BLGT~nsJeHEeI+ySDh@c zp$cv*c8i4^6~Ly1dD#a) zJYRqH@=L9e7^y11%8X-l3tXUaG}dCw81oPrW1AM`&5>hjy%w9eolmj%&i*iR-Z0s0 zer0Q}Ktc#njX;$_nEw^H=dS`-p>?*F+c(2Bo}l>|XaxveWgD)1B7lpeN?_JX^f$>_ z(qMF#50!TzP@4F;ieKHvY{w$=3=! zl#+FxzA|wvT&rHhlwQfBs-&{*1E7jY^;u3PWW~UVR*Mju`O(>Zg1r%L9_>oJo7=$s z>w|s&M+ZPJ=v|>(u-e-6zOx397ixqJ$@a;b)kH&hToa|nIT@7AI^~>I| zD1zmxFX-3*Fdo&raA$4jHRwYn< zz-Hf7u}3r81MiI5-vcgi@8~|S5D=itqETkx*d)XAOzG?5)Don~wbl=f@}$y#cJ4ty zw1ouVdI_FJ@yy+py?1~hI!Z+0H$L9eXVP2BXOt1cqnx>km|n#h5wTV2r{cis)OCEq z*h^{z9}`r`BCaFv;9AeyEV!5I-@E_ZzuvRsL*HabLkN@(XH`DmpARiO&Cq6M6A)9?pRfmz27D)}0NWDcAY`T-6Iz+RO=FU$>=o``7#QdH1a@Pj|4-_gK09Nnue^YXS$L zXxJ=uh%tI-XG>`R{nMKBK@_wFUey^Gx86%|`0j3#1rR}X(l-pQoLXEo`^3eKjQw1I zR@={O_*dX6liNjBn?Y7!W!Yg~Cy;M(_KEo$Q6}GmSOkPoEi?W`xE}T^P zMSkc?3+O8JVLHB3{mxUj$akV~!DVv5#A!Q>+095}B7(-d3U#pS=XINSch!1rd)yuP zcM&}#SA~TY5lDlCe1cPERyGOE2ztw!R&L$O8anoI@poEzJHi8dj={gY7Wx*K7negZ zg%;KlHTc2fn%tlmicWF-6=1EPG-y%M1vui&(gZ9H;Yo=mt01|lH93jGJfZUdItS~) zHrYg%99z<#>D+BDQ1U5eH208sV}D{Z7}IkHEo8Bny%TZA8Q!L&9#V4)jKugNXk41y zmwaQ~Q&5=b>9l$uCX3q^P-l`GXo5Kpe=)5QP^R&j`<+hl^uqgA3k*I%M~eoM zKBur~7za?$*e_8+u9)+YsXxyOl-ob<({Dq4>-#mQu40^(k^)x9@HNMk89Na7S%m=E zL{lI%A6z=}80mqE#w(ugu%tmHAiFDagoO)4Jz{y!n25@4;akg6 zAomW@+_-V=-$DQVuB)b#gZf)_#j2nLgZjJbN|V@|Q0}#EQRekc1_E5~qyTnkpcCAV zjw{+qN-Ho`$T3yaFf&y!v&q>yBgiu`Gkx3!UTV;wg1-vKEU|qUe-*?B1p&bZ0Wx*r z0(c_@Y`Tfz1J0h}ieRRpV!#vDD5YC1mz==1^OZnJi6^~ZVro-t&Z;vqPU4v!hEWMH zfcB)j&gUqkvnJb?xgZp|`vIp=a*PIaft-C4I;{skfi6C{(6866vKyeX3twk^B%Dq| z#t@rK>$NniqQeR>?$t^&g7=uK{I!6rb4z-FD&h;hFxe9_ZPz}3Vl49{k#_a6@ zHJ&zG-vK{te{5yW>gNhk5T{7Cp|M?w=|N^>ROTr>@MdPM2I5uRRQ|A}!x`Kp;vwMv z`Ty+T?*#)%wdf#zcdF&Jg23K-z!jbSKhS{g%O9wE$YB3(ug*9I^v|f2DIv_Cp0y1j z_@7>sJ=fpfn9~&GA1v47;jcuN*9`O@%q?KzFD4yS{8#hiw*u5Zqxq@lf3f6@^}ju@ z>{^IFqih9Wf9K#SKKnb$Sf=$?Uu5~Gzr5+H@xNGpjpAS4b6pHL!tdo;rpm{PydJa= zFtv&6|Kl7?ivMwr>n0M2-xaYyoMx7POJjlR&8+gjo$gf=$ZG-X+Y91v3j%@+@+^?@ z`pL%Llu=F+6bv2wPnqrimU#D%L_2W38TU`o>~k|N=;2!m_r3-G55(7k3q}A6oNUJc zcC^3(U0cxp;E}I5DA-?|_YVfziVNlf@rD&ZAR+&66M*us@bv%N+%fy35DaJ#5Nkvb5XygqI?j; z{|Mm_zX|;khX!V}qrRE7+9CC==hgl=*w;Gbf5`SMuVn5w(m(gj|B$|NylIp{{C)k5 z`$qnsWc+t7{tr1>@Qs|%LH0)cPgV{BLi!IeR_2X()`9oufDk)z!5o#}F#gUDZ^AK# z5_mW;ARv-RARr|F2z#r(5j$F7fO-m0|Ft{+#>)Rl)_k4QeT|$Si2pk#{=Xzsz1Ijy z`Sw%-HM(&5E%pESx~aXsk)<)Cp{4mhZvR_Z)WCH-&+Ex6zMlMt|LA)4r4)Ujt`^Hb zCjJ&+)-{gIe`FA&EAM|eNe#1Zqh#rNP2L=KrCKyeBg#L7SOimuVSB9Jl^t+-`+$0hkNvj zH8}o{4yIoGKe${kF4&UG8y?w9`v+@%#jriz^xiwc18oQY>P6|p1v~ZmtAp&#U`q>vew*UL8`I->+ zI)l<1DIl)~%M=v^r0>-o|9y&tBVS3Rz@7d-Q5g8xj|+C51&kb~2J#Q!{8@(904~^% z+`r~R0ka0M|1R$p?=N^gKpbFhHu--YA~3B5qvh2y;#bS)|F2oEZOvDo($^)C0{__) zpFv!(!0-PvN$=p_vpIUj-Kt)3d|+cX(tk}NIElG?2o3_``5I1n{+VR4`b}9!IUX=| z=O#F@XKi@L1gn#4QYCABii3Er#geU>Ymx2ipYnlPJ=;LRf2&<1^^Hd0i?=%vWe)lXur!|=%y7h005ww zD6WB*e(rka3I+HdWGN-=qs9dNQ+dPwSb~_LHXbYW7|^cJPoi+xuTLT;xcnzkJw75f zJoayWv+CehV`u;X6CMCy1%PUM@IWuoNq|YqjfUc&L-V%TAQ`q6~*tj&`sguwhTO%uU4!h#r3Mmu!awB zMcU^xs;s(*0uL8@gvw{FxPMm7iPb@;=77xU1U;mn8;dCTi*R&1IvoPhD|))o$F~Ya@V9640_+>KPnA_$a!~8+ZL7fe6RJ8}K!}{o$kAyz*NcLdo zOV2?B@`WlThu9kLuTOc%#JHLa>KY&DFqncHAAa)ie=)^D$mVV%k-&RUf_nfOBdIqK z$ubJ2`LLc1^rxmx!Q>5F6P&uo`h?tPoWZ-r6pD0`*u;}a6>MS^ia1z!xD0a#{w-!ufe$7bZzB=qW$1K|7S)^n zP6S45P0(60K3Nna$$ktMVTC|{B9VA+yS0T}5hy}~HDFzQ%gKQhB7 z`Xm&WEM`NgD~f+klYKz>JD$}7$ZWpA0susi$-n`q$H$_Adx}{h{P;gHJ4vX8fc`sX z^O{9o5nuoS4pbmcN(La3KkmHGgXj@w`5dLOgnWzueu8*&zXJ-9K7l2Uf^k?k3cN7* zZ9z|Y_SInptqHv+Hukd1)-xUv=OsFiZ7`i92{ujoaN4Eo?Bxs4n&a%}sO2JMV8=5UT6y`t!Y zpuCB%MxK-3o37hr%4D$D&LJz&J)JyI2}n?I0$xYhN4i5G>MWiXgJ+a}c=rCGV!G?y zOc%I|@$v-f4kP@YTi)Z+NTRIZA?4Qt()7dVPayIyV8U~7rbOBMlQCx)7@Wzp9=KH4 z1sWm4@Bri=Fcm;peoB%;iLx*6Q7t62ebRfsNZ1zg!L*ZdS7W9n*CE6C;8kn4NN(~% zzalCS`p6{jCs1*^MA?XP*5NM~G`|@-KZdF9EtMeXli!<2>2Wa#a6&%)3UESl<(+?b zuDu2y-15BJgKHk85uS(`&k-4H;+T?G7_0BohAM=+MpOWlpd<1`+HefNW@zV=r-yw$ z=rtPNCcV={d9v+!Ymafh+HNFzlBm_zd(9Qz=7;W%`kh>4mW+(w_t}zwH;-@0!}#Xg ze1vq#4y#7F#+0xWQU{_MRvLlUc$-(L=`WV`z@kO8T)sg zY0%CkRWSz-=HEflJ)#{6d)E1=^krdG_c?0%~ zFKl-Oj?Jj#QPAcVlG2B4K^iSM%N;1q>7T5*kcg4cMnj=JSJgkx#Uup1rxx}g(auQ~ zP}7P`1>(pmcIHX8MCK&f*@U;HX~%u@q^hGSLRxM5JSkVtZ;fBtP^Zee>g&6}rC#-B zt;Iy50q+kdoM+~En9O-^XhJ{vN z7(IdOiSnKTsq?$i{_YR%-Gs`vo`YZlW;3cgYT&q5p+}pQI>qDy^_0mfF(24wMj*?W zs>w_e8$s*#*W~!{8^`@KnpS+T{O1ls6R3ykYT|c6HgR@MM*we}H(VPEXjEpngacz;_KI~_;e&-Oq?T%3EgdSn9Jk~2-)TS=V({XS`up9t;e zRt(xr&naP8y*vL>8E7D>jh=1TKYd34rq_w(+hX|O>Au2p#^`e23~yBvWti`dv?=nS zyoG(jbfb&QHffR++z-hms{Ocwq-xwAJ9%J&4!-fRg5}_qzoXjZm#qd@&Z~-YHj3vO z8Z+ilzrG5jBY(ff+a0}X-VrvV-{m~=&HxU(NOtE)e?5OD-6b(G6X0QI$E@KrbD23d zr=&tb3##O@Pc<1LEj^Mf^Ddf7WXfo}fD??gj-<3T68z8$ozvHN@wvq@J;Alb6R8SV zRJb9?>E+XGu-0rYh3o5wK9@s#L6wRy9?M(NMqX~YnF05rFy^vdQO8=ZFt^x9b>HmL zC7HY^Oc6d?Yv0CKwotBxgeUbfEKx~>DgB+gO_r)QMu|TKg(g4KMVn?I-*wr_P=~kD ziorYE79ymwz+{XwlT$;6jNa>V|1eQtyr$r8vmBSgGU5CMWBh*J57L#y-d2J|ZO|Hi zRHmpy-qyOI78wXk01Z0}Kg6H=1Mbyg|&KeU8 z>28-8XxRfzqcHxOgcO>$gcMeS9vXCY4YQOcvV-eH2uxz#oI7&@Cr33b4vGdq-7h$o zKFKf7Lz!}u=4?ds!MoCSJybO5TJ1G=nCklzOkiQG)lWgN!otrTKiGPIKVGY>8Q*2xcnO>9J%AMRU zn!uO`SB9(5ugEn?d6S*CkLJQ1n9M%Sq#xURSgA1oY#?Nxbm6CDc^#`U$y(=Ut~syt zbmei=grlEo#JJP*_5wUU{DuRpXUzS@VX`Cq;Uv{jA>p?N@ z$s|+L)Nigqh+2t$S;Vch4}8C-W`wQRU3yV9#x11tD>P`VnPJAuVEXHU_`7fObyEXx zKb(jT?I$osa^o8&cq$Pv$UpnS-xy)gbS8kHI$r(~Jr+x*;sAp&~AcT;Ag)reQ>QRWj$B*l7ymbBHc$mnYdzLCU8*1 z2cwddYTKj9C3@zUrjZjkOz_HJza*2g;ZCgeNIeU3_*$lRQ&g?wgP$41Fm&Z}ayAN$ zko(Y7#dNPuu&6UZX7K|kx`?(Rymr9jb)$nb3zZ`Xt%s%6nSp}YGsa*orIExB@yiht(*xMe)a{gFXAZlX;6eAGy5yi&an?17m zwT;o%6%@AnqP>LHj!|>4-Sfh-v!Ruc%}0dx#msRTb@O|MGx|~idLbF7PNXr+4u+c) zCh1HD6%*=nlTS8bX|iWpl4iI`kwTxTd`vi@>;lF^M84d-pim#3KAz<3xvm#uHM(o-Jj%>V z?Swi_w>6~^0rj}ibCCfiPTQ!J8AWqOdC}}UAN66<-`H|~25JcK2KC80uN*oztq+QMtVsN~h)%IoYpvbB$LJSC+luBv5xt6H#zoc3@gT7}J|G+nDd#0sM z8ujCV#Bz9f-97m;2XT4XtFqYevi5lf(hDB+c1b&{y|)MauMuE?f`09Ro3*=44h(q} z(mBup>+xrZtiWQZJx9!Q(tI}G_u}ind`j=gIu~CB5U#0(WFqXpxCZ$lBL=eN9b`ILulBVSq zJ$G6VQl0=hbk;a!=;ai%I?F$>G)lf&mJ%wmp$B-GWE8E=Htb@Pr4@G6zG~&BJN24S zT}C~!>kIcZXlYZ%x;lz)gflE2t-)>`x(atTHqJN}3EUrS?i=NSyaa#XWm{xp8C$V< z37kF9)RF($ExF}N*oqM|me3JmSdd$&Vs(%0VI>9p!lPwM8u}g0`p`r#e4|2K_gK2l zTB>zVx{NSnkPh(cS;IYxcCKtor~FI6ve_ODMKQbaFWRt4ys!zD_f@k$_RwIby^M)! z9im5_%hN?=#Lz3nqHxb}E~&3j2^w_2BI&aP{Trw5mA-Sp^j_s>Z(=uc_)EgkkYm~nqmci3gDtuAAJ z%P#6nno64bk+fUrzlA8klS6&GXoeQjhJ0?B6}DwO4AxGhHgr{YjQVMwRCqz%g1f%5 zywNYkU34!xb$=CC-65X&jqKXipbO@nS1LyYxQ2Xkv2)SxzD83EJuT(-!ahP05v!Dv zP7Ai*HcS#NB#%6J`ppTM^HwbAqyK`=9)SsfUAL@g3=sq4*uIdxm$LVKt-bX-c)g;Om4y8yld4ub=91rR9`Jf$kn3()nO`g5hAGI->j7dDBLoUNy${5{B zz%EWD&X${~h4y+PK>;a@rrU!)!$m>7?^Dw@F4x?t*69Z~R=uNFN_?V>(iBM(LG{Q8 zRYnog#j2bv&Q=o4Tw}m-(JyJcL{^@S9sq^V4Oxkn3FH<9!ibUxJDy&~;|QS@6pieV z#7RjdIv#VAK+-#(-kBjhV_zb{>`D4Opso2QB44;xk??o#lV zpe~DCJ8*8K*au$+c8_n*?iKqnR4hWh`g+i4+4!7hSEzHP6vtM`w38(g)205aY=+Ba zb%&sHuRX7N27QjPbWKN~OIR7o26g|Jb*X3Uori55V_orK8-|W!5uNL$qXSNDc}G)Y zTV57WIZj{-+&+y2Hb;yp3NVF)#gUlNN7~dFYG1!E(S3M&$MO3wqd7+gT-L)B4)050 z9;}8}lHimu$=79yWW6@JXE6R}FR4C=8wKuhH>gL3a#RZP54;f{0AToMHwXcgJAe*M z*7Y(bnIn40rje_2b$MfNwHbs6B~)j@tt2mLEeu^`sH>tL*|h(vb%DSc5W-_be>twmgEc^cet&9>?{t$Z8cWwl>m06vXeZ@XUO1Fn0&3_Kh`AcuCY z!baZy%wdSvQ=VVaW7pfl8m&hfE7=3Iv)J`8X%`LQ3p*QRa2z62`?NX9fWKi(=oD>= zh5kA5e1fZ8CZc)t^!a9|Qbb$$&Jk9-)K{<3$+n#XRoIi1TM@2yzL4vYpMgZ8m^}y< z_X((xeJ%7M7%`Lah$kIt?ptqkd0CCP&pnJ_uVV%dMF0(fHO@)sD)b5m})xkz6J4V}J9JpafcN+HAENvYUL zIFNd3!Fj~2=8tC-vvqJI3%jma(_YNZD+r!%v}ka$b{;)UH}WroG@(lYLXupKT-zDG{?CR5f|Cl1$$4z(WI;3jRF08)%5CRK)t5wDdt zzRFz@cggPsJ$sd+0&=N4tU!_KSHAIT#>j68WQOs_RGk$Mh;<}XnuVHUFBNK#EoT%{ zw_l4!rJyo|ov{wCgBxRMz$>N}>2L84$g>q^Gb*i`Y)}ozpfbVs`pt?e0`m4b%S9or zFJ9>eeE!7M*D38Fz{msF{3g+?aDPct0mtXgbzIHMFMx2}@)N!Lay$^BXG28>pUK)T zK#-B-Gud#|s#c{Anf|pz(p>?Dd3~_e>mr+3QC6|{QF0ly3{jy}6B(4$Xi$RIVVlvK z)Be=q$izc1I~2vTFv=#fE#|428kas$n-|nQ-npy36)tumhguz?^-iidjCL?LD%$F@ z>!o=0Zpw|G%|d*f*UG>#at>rTa58IbQI@nVTX%ZmS&remM^?5*Ng?4#za>x6m2aMYMA7g z^meo1iS@iJW!Yw~W2<~$enUjAbLP5IvbFo6D#TQ`o?h#xBtsx$ zxwSDJG#WGV^;;9b;E&R$#6{8V32c=rpQ=72PRQql$7p+oc@%E3-!QiYSm9wiiZCen z?yRqs>B9}>)84NhR}5aR*}U0MMaWY-;}aG|khLh*Ivqy32ybqj;9~FRF2pCFchi(9 z!=Ktq`R=XMSUcyJQ^@MVnp!H%sJ{%xqR9S);_DxuI59W`ylN?Ynch}fW1R8v?Txdx z)x!2bh*ierb3-mNt1r)tauK&rLe|7ISEz5;uLwR=Zz}OWlx`W_co9CRof(x@1T5eU zf|%}y#2KQ?J~3S2+7t6^tIWR__-_8wHXygWDPexr)G-1c)^d{6lch@ImqNml9K-b; zxF*k?0`1mw3MhPEZdtWido*R2?9!<*0B7~F<3q}8b@djV%8(cs-_nOk=NPe;^yhF_ zPGv@#C7sU+nzco7!^-l$xh8K}6x*R`&Z#0&SK(7x1hQ1NkaNg41bmm^bR@yfUYC>W}=QJK0Q z(PmQVrNnuGqe8WXs&oJI-4E@XHgMNz&Dv(u^%^|_fu*aUjAU>%*?`6jrCj+tByjRei3g#tG$ARkjH`z6;HUgICF&SNjI>H=idAiv16?-i%m` zousWaPmVLsr~BZz8Omo*^|+0k30J{$<{ETIzhW?}FCkLu|2jd*6SGuHH13osaZ^L% zg)nrx1EqKB|5nL>&H63+EjDWfbZn|Q38IkZNIwsd+RJ#6<=dd?%yg(O$8+|?R*Y|` z_@XSnIKkh(amJgGgg4Aipyw_X>}L3Xw56J&)TyjXA#GgwD8w1+8*L~(Q%u^kCH(`d zAKO80IzgUX!i*%YmF&mz(utJDmLd6+@d!_fU~a7im$KO{s`iU*O|R<3RsH?L?61e! z$HxQNzz*(}4z1(0^vShvQr3n9uEJu=SMAEeM&ub1V|~}tTlG81eKS`dQHqCwuQ7XV zS$ch}`)^TC;rP)AlD<+ut4T2;$j3N9yzn&$Mce6mht1YAvGAs?W~65yn%bE=MRG`G zlo~&+yC8I772PxXf;LKh+aD?V;%0)6NDj#`0_gABl+n`g9jAgf-t;4g^7?}bG14)I zL7q(mg(LG2%UA1zyzho3EwTFxyD`hQQ_*@wXLG~gXvQkZEH`03rPgZq?*Iutd?l6q zZ3g_Q51^6stR4Yw{%ej(!D>M*_)^z5>HJ$Iny;~P6uUn7E6m`G+iVKIznt3g|2hH} zwEztN#1UuJ&$gBDrt(UgK&adUul#Dw7ADfhKoX0DaSI&$eSzC<^%=cpX{{}PhH7~H=t6h?%qv*M*qaD z)Yt|d5wjtRUY!XY4i>rTHimLM+O3U zFB(O^E@5AwM3u!Tq{u39M>V!BRqGNcHnCbQi5=vW9r|(BuEY;`E0h8XNxzP=4PN|u zn?lHF_^qEdB6U;nT-N6M6)deLwD{IkE=mr}K{gn+H|lw~a^X_=fj=HT8zbsAJpCtL z6`%_o(jd06jWs8}^N{g5gDog>S)j>+LWKjtZ?t#Sd9C|&WmssNt|?!CI!K43pgH9{ z?t>!o%PZW`dF7u;JW;hsXo8>D)z2*#+VC>6ix!Lycn4||G+5UvRljUGh!SKf5oe0T zAy;DG>4?Jp{S~ObMCv>8<2RnIDra*H$T!JLR8TxDG3c}(9;7^o2+Y>r*TK=jeGuV# z#u-cO!i<}&y%b@L7`=pSEhAnitBI*vUMa;DF=IlP{X*Q~$kfh`pv`F^pfQuxSFzvs z?NM{iY<1wF(G2kYQAawg{F}?^EYOd&wu6&ZzjeQhq4lm6*cyDD-8}ps>6cvlX5bSk zOax*6y-<=!VebV5b|Cuc7RS_E?=!^sc&=RrobV91A=Ui$LL&nKeNSO+wF(#Jr-Y(!64PlX^zHQPL|;AX#fc938wysdUIC$F+?bU}WU=ZGnwk0Ij^Xjwi1)qzLS=0?9G-1YrAJ(>G4R_HzT^kTWX9C zY1${Fyk4fKY)t1|8tHvwS@OCx4{C0y8ufH-l%{=7bO7v~ZJ~&&1%P%fHYw50ta=7w za?zd)DY0Y?P6ZAVIjm8Xu#Trskt@Ei{DimNpKQ$wtwbuLm_>fvprur+95WY*dk|(T z@a+OkD~@(~0!xZ7IW34^df1*DOD$3>+g5z|=j@7f34ZJVm$JU`t9$g9w61LxY#cTt zk?rAX?pd*ANL4oA(0bVC&Frr^w>Cv!hUOULGUoFK&<21iuXY+eih3LP)n?xmu_*B# zI+IvL2jikFn-x`*e6u&4{Mwu|R}w}htH}q-z*?_d3bh;$wXzV7bmPkI^5FVR`7A$2 zfe?+0_)OR*ja+;XWr3x}AK3zfw+d@L1$>!iJHj7NV9_3yTB>;U6^W?&v=xiMygho} zynS*9?P~`xsu@=C04LgF%sX@OWaD*oN;!@|rRB>LpVQ7*HG_9&P@WjZ&hR?FZzRP!Elf4PnXa|f@tCuu|HSf6c2-d4xU=v2U9RDFLyCM&?5Kf5 z=Cd29RFNCC!Iy*}m++2;Q_*ruJ6+q%UQ?usv_`!!pxvl>oe(UO6Gh~yEUQ-SoTn{U zu(C|NM;CWf`XDE3&Ib)~!%QnuYMmgAi=QpFCt8FlvsC1IUW`(0<NEEWvZZrd}^{XhO7$XqVmQM!O3^-##ayPlkupdr;bZW7J7!=ojXfgHOezClz zl3Ov(QC&S1zg%t|cgF+uqs4K8kegM_?Q?qLZ){hCbj9V+w-|DPvP7w=8E=Ulx)^Ip z-cz)I1s38S(Av7azaDi^=1SEz9v+D6_`YwtDF>>Ff|Q2LR24i+=|{rN%_B(1uTZ&p z2-0T1w*<#!K3wBWEWRtY@*BMv6zJO$>iKBe`U%T2FQd~og^?u zxvVeeUfYDsgKP!;5gL!39Fn3<1?uNGDcJeg{AYy3k<7YOaZ<>JQu}yGDaJHve2YQt zYcGMCkhh96Kt_)dy8DkLd>@v5dZ43eb{o`n&&CvkLZ2JVTIf8-)T>)C19NIr0p|m8 z06*4z@?fu8MTt_MShPhy=D0g-44tazSKT}IUYge-J4*(<(cATRM#B4D%29c z7rq)ew|GXE@@~)-thkLXXV}v|?V;I<>6!&xUCfo@4lZNEuEexXVk$V$W!y_#9SbJDj(wAA-^byyw_mAj zL&Zuvnm*uRH%m2XL2n^8C2DTel*s58K{;(Q3D%IYowp+j$upT)HI0zIqQ4Dz!;kIa zd8?0P)^EdS>zHUW05&ELu6^a($b(|c2S%G&CHlp7I1=ku%T9A3i$M!>nyAgn)dht5 z`$=L;X`9dh1psKn008^|Q1GB&{is$a)KgADtd1rk8VMwmU@jIGs0PWic?mi)5+z4} zsjPO&G@-^(qsmd!z{Od^#jEA)iH5Gm!NtMz1tf3Gh=v4Ff*$O*nt&(~g942yfV?l; z9AT(153hBo^Wj3M;FMeCAv^dsG~}7wpP}eCRBnG}nfbM)XD_B-*=20oOyp#+He%i| zpQUqMZJwLLxrP&?+1eBY_6AUX%9{0#!8a-pWulOLC#{DVt>~&LoAhBg_slw$Tf&=v zUNm~loW>zbd14@4*8!*CUTpLO<3(Jmz~(B1(&UX!PGV+bfj05+Nx>p>Uv_YS#K!*0oR2dUU%a zs3&ZJJvri2ZP1+@=Bc*lCHSXC=Ld#+lKl+2{iAUk@(JNdHWNMlM^+eH_GB@qFzcu%Gbvd8g1=zx|S%u{=zI_>~0*0C0oys+sB;3rqgEHdovT z`viYm0sY5~chx63EY2y5fC520*I6#y?)e3@Zb4v|Dqd$$%~C{yVte3IEhWN!%}T}KVNCxi!8{W11$q7nc=`L8WNL+~yIq`jn|@BP?M z1gHbV(8BDGq#^@!kDmU$7yGYvL$DIhA5VYp%o0BlbU+9`i9X7Y2WTD%{(Xi1-`<{7 zfLP~v1pU2E{XJxgC6GY5hWu3!{>wnC(IZgRf)w=Y9|E)yV(50$M|sdF0_c@B zKInKH^^s_g2JKu9;-wOVjPoy|5(|g|B+f@@9)V&6)(=An06-*U;)whO6t#L}#v5gP z1d`g}*G@p#mLab3{RK33`s2bAbg{tnXM zsz>JeDTYTN2q2XHySslvSOxxK9jF6QOv^p_ZaYm3o!0QE$kvMeJ2IZ+qY(K?)1y3L f0v?1uL-k1d_p|9=)`sBlR>+fz6CSDUxA6Y}2O&Uf diff --git a/gradle/wrapper/gradle-wrapper.properties b/gradle/wrapper/gradle-wrapper.properties index 0f3a6fe3..ab8d9dbe 100644 --- a/gradle/wrapper/gradle-wrapper.properties +++ b/gradle/wrapper/gradle-wrapper.properties @@ -1,6 +1,6 @@ -#Wed Sep 02 11:48:49 EDT 2015 +#Mon Mar 07 20:47:12 EST 2016 distributionBase=GRADLE_USER_HOME distributionPath=wrapper/dists zipStoreBase=GRADLE_USER_HOME zipStorePath=wrapper/dists -distributionUrl=http\://services.gradle.org/distributions/gradle-2.5-all.zip +distributionUrl=https\://services.gradle.org/distributions/gradle-2.11-bin.zip diff --git a/gradlew.bat b/gradlew.bat index aec99730..72d362da 100644 --- a/gradlew.bat +++ b/gradlew.bat @@ -46,7 +46,7 @@ echo location of your Java installation. goto fail :init -@rem Get command-line arguments, handling Windowz variants +@rem Get command-line arguments, handling Windows variants if not "%OS%" == "Windows_NT" goto win9xME_args if "%@eval[2+2]" == "4" goto 4NT_args diff --git a/spring-kafka-test/src/main/java/org/springframework/kafka/core/BrokerAddress.java b/spring-kafka-test/src/main/java/org/springframework/kafka/core/BrokerAddress.java index 72e7a571..0df06b74 100644 --- a/spring-kafka-test/src/main/java/org/springframework/kafka/core/BrokerAddress.java +++ b/spring-kafka-test/src/main/java/org/springframework/kafka/core/BrokerAddress.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -14,7 +14,6 @@ * limitations under the License. */ - package org.springframework.kafka.core; import org.springframework.util.Assert; diff --git a/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaEmbedded.java b/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaEmbedded.java index fa7d5082..8ae15fbc 100644 --- a/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaEmbedded.java +++ b/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaEmbedded.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -14,7 +14,6 @@ * limitations under the License. */ - package org.springframework.kafka.rule; import java.io.File; @@ -116,10 +115,10 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { int zkSessionTimeout = 6000; this.zkConnect = "127.0.0.1:" + this.zookeeper.port(); - zookeeperClient = new ZkClient(zkConnect, zkSessionTimeout, zkConnectionTimeout, + this.zookeeperClient = new ZkClient(this.zkConnect, zkSessionTimeout, zkConnectionTimeout, ZKStringSerializer$.MODULE$); - kafkaServers = new ArrayList(); - for (int i = 0; i < count; i++) { + this.kafkaServers = new ArrayList<>(); + for (int i = 0; i < this.count; i++) { ServerSocket ss = ServerSocketFactory.getDefault().createServerSocket(0); int randomPort = ss.getLocalPort(); ss.close(); @@ -128,22 +127,22 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { scala.Option.apply(null), scala.Option.apply(null), true, false, 0, false, 0, false, 0); - brokerConfigProperties.setProperty("replica.socket.timeout.ms","1000"); - brokerConfigProperties.setProperty("controller.socket.timeout.ms","1000"); - brokerConfigProperties.setProperty("offsets.topic.replication.factor","1"); + brokerConfigProperties.setProperty("replica.socket.timeout.ms", "1000"); + brokerConfigProperties.setProperty("controller.socket.timeout.ms", "1000"); + brokerConfigProperties.setProperty("offsets.topic.replication.factor", "1"); KafkaServer server = TestUtils.createServer(new KafkaConfig(brokerConfigProperties), SystemTime$.MODULE$); - kafkaServers.add(server); + this.kafkaServers.add(server); } ZkUtils zkUtils = new ZkUtils(getZkClient(), null, false); Properties props = new Properties(); - for (String topic : topics) { + for (String topic : this.topics) { AdminUtils.createTopic(zkUtils, topic, this.partitionsPerTopic, this.count, props); } } @Override protected void after() { - for (KafkaServer kafkaServer : kafkaServers) { + for (KafkaServer kafkaServer : this.kafkaServers) { try { if (kafkaServer.brokerState().currentState() != (NotRunning.state())) { kafkaServer.shutdown(); @@ -161,13 +160,13 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { } } try { - zookeeperClient.close(); + this.zookeeperClient.close(); } catch (ZkInterruptedException e) { // do nothing } try { - zookeeper.shutdown(); + this.zookeeper.shutdown(); } catch (Exception e) { // do nothing @@ -176,25 +175,25 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { @Override public List getKafkaServers() { - return kafkaServers; + return this.kafkaServers; } public KafkaServer getKafkaServer(int id) { - return kafkaServers.get(id); + return this.kafkaServers.get(id); } public EmbeddedZookeeper getZookeeper() { - return zookeeper; + return this.zookeeper; } @Override public ZkClient getZkClient() { - return zookeeperClient; + return this.zookeeperClient; } @Override public String getZookeeperConnectionString() { - return zkConnect; + return this.zkConnect; } public BrokerAddress getBrokerAddress(int i) { @@ -231,11 +230,11 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { } public void startZookeeper() { - zookeeper = new EmbeddedZookeeper(); + this.zookeeper = new EmbeddedZookeeper(); } public void bounce(int index, boolean waitForPropagation) { - kafkaServers.get(index).shutdown(); + this.kafkaServers.get(index).shutdown(); if (waitForPropagation) { long initialTime = System.currentTimeMillis(); boolean canExit = false; @@ -255,7 +254,8 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { if (Errors.forCode(topicMetadata.errorCode()).exception() == null) { for (PartitionMetadata partitionMetadata : JavaConversions.asJavaCollection(topicMetadata.partitionsMetadata())) { - Collection inSyncReplicas = JavaConversions.asJavaCollection(partitionMetadata.isr()); + Collection inSyncReplicas = + JavaConversions.asJavaCollection(partitionMetadata.isr()); for (BrokerEndPoint broker : inSyncReplicas) { if (broker.id() == index) { canExit = false; @@ -279,7 +279,7 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { // retry restarting repeatedly, first attempts may fail SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(10, - Collections.,Boolean>singletonMap(Exception.class, true)); + Collections., Boolean>singletonMap(Exception.class, true)); ExponentialBackOffPolicy backOffPolicy = new ExponentialBackOffPolicy(); backOffPolicy.setInitialInterval(100); @@ -292,9 +292,10 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { retryTemplate.execute(new RetryCallback() { + @Override public Void doWithRetry(RetryContext context) throws Exception { - kafkaServers.get(index).startup(); + KafkaEmbedded.this.kafkaServers.get(index).startup(); return null; } }); diff --git a/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaRule.java b/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaRule.java index a6c98b4c..b7eeecfc 100644 --- a/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaRule.java +++ b/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaRule.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -14,7 +14,6 @@ * limitations under the License. */ - package org.springframework.kafka.rule; import java.util.List; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java index 102d0049..5fe2199c 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java @@ -343,7 +343,8 @@ public class KafkaListenerAnnotationBeanPostProcessor } catch (NoSuchBeanDefinitionException ex) { throw new BeanInitializationException("Could not register Kafka listener endpoint on [" + - adminTarget + "] for bean " + beanName + ", no " + KafkaListenerContainerFactory.class.getSimpleName() + " with id '" + + adminTarget + "] for bean " + beanName + ", no " + + KafkaListenerContainerFactory.class.getSimpleName() + " with id '" + containerFactoryBeanName + "' was found in the application context", ex); } } @@ -356,7 +357,7 @@ public class KafkaListenerAnnotationBeanPostProcessor return resolve(KafkaListener.id()); } else { - return "org.springframework.kafka.KafkaListenerEndpointContainer#" + counter.getAndIncrement(); + return "org.springframework.kafka.KafkaListenerEndpointContainer#" + this.counter.getAndIncrement(); } } @@ -414,17 +415,16 @@ public class KafkaListenerAnnotationBeanPostProcessor @SuppressWarnings("unchecked") private void resolveAsString(Object resolvedValue, List result) { - Object resolvedValueToUse = resolvedValue; if (resolvedValue instanceof String[]) { for (Object object : (String[]) resolvedValue) { resolveAsString(object, result); } } - if (resolvedValueToUse instanceof String) { - result.add((String) resolvedValueToUse); + if (resolvedValue instanceof String) { + result.add((String) resolvedValue); } - else if (resolvedValueToUse instanceof Iterable) { - for (Object object : (Iterable) resolvedValueToUse) { + else if (resolvedValue instanceof Iterable) { + for (Object object : (Iterable) resolvedValue) { resolveAsString(object, result); } } @@ -484,7 +484,7 @@ public class KafkaListenerAnnotationBeanPostProcessor private MessageHandlerMethodFactory createDefaultMessageHandlerMethodFactory() { DefaultMessageHandlerMethodFactory defaultFactory = new DefaultMessageHandlerMethodFactory(); - defaultFactory.setBeanFactory(beanFactory); + defaultFactory.setBeanFactory(KafkaListenerAnnotationBeanPostProcessor.this.beanFactory); defaultFactory.afterPropertiesSet(); return defaultFactory; } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListeners.java b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListeners.java index 9d835db3..b9f97199 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListeners.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListeners.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.annotation; import java.lang.annotation.Documented; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/annotation/TopicPartition.java b/spring-kafka/src/main/java/org/springframework/kafka/annotation/TopicPartition.java index 49eff4a7..c2472586 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/annotation/TopicPartition.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/annotation/TopicPartition.java @@ -13,11 +13,11 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.annotation; -import static java.lang.annotation.RetentionPolicy.RUNTIME; - import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; import java.lang.annotation.Target; /** @@ -27,7 +27,7 @@ import java.lang.annotation.Target; * */ @Target({}) -@Retention(RUNTIME) +@Retention(RetentionPolicy.RUNTIME) public @interface TopicPartition { String topic() default ""; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java index 9ae4d0d7..ebe0ce2b 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java @@ -60,7 +60,7 @@ public abstract class AbstractKafkaListenerContainerFactory getConsumerFactory() { - return consumerFactory; + return this.consumerFactory; } /** diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/ConsumerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/core/ConsumerFactory.java index 59f07856..8007506a 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/ConsumerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/ConsumerFactory.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.core; import org.apache.kafka.clients.consumer.Consumer; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java index e3df8b94..e4162f35 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.core; import java.util.HashMap; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java index 7988205d..297ad3d3 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.core; import java.util.HashMap; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaException.java b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaException.java index d578be11..52d0e069 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaException.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaException.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.core; /** diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaOperations.java b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaOperations.java index 05e48afe..bc29fdbe 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaOperations.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaOperations.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java index 0b84d29b..523515e5 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -55,7 +55,7 @@ public class KafkaTemplate implements KafkaOperations { * @return the topic. */ public String getDefaultTopic() { - return defaultTopic; + return this.defaultTopic; } /** @@ -113,12 +113,12 @@ public class KafkaTemplate implements KafkaOperations { } } } - if (logger.isTraceEnabled()) { - logger.trace("Sending: " + producerRecord); + if (this.logger.isTraceEnabled()) { + this.logger.trace("Sending: " + producerRecord); } Future future = this.producer.send(producerRecord); - if (logger.isTraceEnabled()) { - logger.trace("Sent: " + producerRecord); + if (this.logger.isTraceEnabled()) { + this.logger.trace("Sent: " + producerRecord); } return future; } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/ProducerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/core/ProducerFactory.java index 92b57dd2..8d9e1889 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/ProducerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/ProducerFactory.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.core; import org.apache.kafka.clients.producer.Producer; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractKafkaListenerEndpoint.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractKafkaListenerEndpoint.java index 73fd9ed5..d629912d 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractKafkaListenerEndpoint.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractKafkaListenerEndpoint.java @@ -1,5 +1,5 @@ /* - * Copyright 2014-2015 the original author or authors. + * Copyright 2014-2016 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -72,15 +72,15 @@ public abstract class AbstractKafkaListenerEndpoint } protected BeanFactory getBeanFactory() { - return beanFactory; + return this.beanFactory; } protected BeanExpressionResolver getResolver() { - return resolver; + return this.resolver; } protected BeanExpressionContext getBeanExpressionContext() { - return expressionContext; + return this.expressionContext; } public void setId(String id) { diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java index f0d9493a..cc54017e 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.listener; import java.util.concurrent.Executor; @@ -127,7 +128,7 @@ public abstract class AbstractMessageListenerContainer } public Object getMessageListener() { - return messageListener; + return this.messageListener; } @Override @@ -159,7 +160,7 @@ public abstract class AbstractMessageListenerContainer * @see #setAckMode(AckMode) */ public AckMode getAckMode() { - return ackMode; + return this.ackMode; } /** @@ -175,7 +176,7 @@ public abstract class AbstractMessageListenerContainer * @see #setPollTimeout(long) */ public long getPollTimeout() { - return pollTimeout; + return this.pollTimeout; } /** diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/AcknowledgingMessageListener.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/AcknowledgingMessageListener.java index f9f6dff7..0cae85f3 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/AcknowledgingMessageListener.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/AcknowledgingMessageListener.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/Acknowledgment.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/Acknowledgment.java index b6aa5705..565fcce1 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/Acknowledgment.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/Acknowledgment.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java index 07ae91c0..3b014fbd 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -65,12 +65,14 @@ public class ConcurrentMessageListenerContainer extends AbstractMessageLis * @param consumerFactory the consumer factory. * @param topicPartitions the topics/partitions; duplicates are eliminated. */ - public ConcurrentMessageListenerContainer(ConsumerFactory consumerFactory, TopicPartition... topicPartitions) { + public ConcurrentMessageListenerContainer(ConsumerFactory consumerFactory, + TopicPartition... topicPartitions) { Assert.notNull(consumerFactory, "A ConsumerFactory must be provided"); Assert.notEmpty(topicPartitions, "A list of partitions must be provided"); Assert.noNullElements(topicPartitions, "The list of partitions cannot contain null elements"); this.consumerFactory = consumerFactory; - this.partitions = new LinkedHashSet<>(Arrays.asList(topicPartitions)).toArray(new TopicPartition[0]); + this.partitions = new LinkedHashSet<>(Arrays.asList(topicPartitions)) + .toArray(new TopicPartition[topicPartitions.length]); this.topics = null; this.topicPattern = null; } @@ -131,7 +133,7 @@ public class ConcurrentMessageListenerContainer extends AbstractMessageLis } public int getConcurrency() { - return concurrency; + return this.concurrency; } /** @@ -207,7 +209,7 @@ public class ConcurrentMessageListenerContainer extends AbstractMessageLis int perContainer = numPartitions / this.concurrency; TopicPartition[] subset; if (i == this.concurrency - 1) { - subset = Arrays.copyOfRange(this.partitions, i * perContainer, partitions.length); + subset = Arrays.copyOfRange(this.partitions, i * perContainer, this.partitions.length); } else { subset = Arrays.copyOfRange(this.partitions, i * perContainer, (i + 1) * perContainer); diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/ErrorHandler.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/ErrorHandler.java index b66c1a07..a383df42 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/ErrorHandler.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/ErrorHandler.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -14,7 +14,6 @@ * limitations under the License. */ - package org.springframework.kafka.listener; import org.apache.kafka.clients.consumer.ConsumerRecord; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistrar.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistrar.java index f850a753..c5186a26 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistrar.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistrar.java @@ -195,7 +195,7 @@ public class KafkaListenerEndpointRegistrar implements BeanFactoryAware, Initial } - private static class KafkaListenerEndpointDescriptor { + private static final class KafkaListenerEndpointDescriptor { private final KafkaListenerEndpoint endpoint; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistry.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistry.java index b8019fc4..25808b3f 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistry.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistry.java @@ -204,7 +204,7 @@ public class KafkaListenerEndpointRegistry implements DisposableBean, SmartLifec ((DisposableBean) listenerContainer).destroy(); } catch (Exception ex) { - logger.warn("Failed to destroy message listener container", ex); + this.logger.warn("Failed to destroy message listener container", ex); } } } @@ -270,7 +270,7 @@ public class KafkaListenerEndpointRegistry implements DisposableBean, SmartLifec } - private static class AggregatingCallback implements Runnable { + private static final class AggregatingCallback implements Runnable { private final AtomicInteger count; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java index 54a90d39..86f8810c 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java @@ -117,7 +117,7 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener * @param topicPartitions the topics/partitions; duplicates are eliminated. */ KafkaMessageListenerContainer(ConsumerFactory consumerFactory, String[] topics, Pattern topicPattern, - TopicPartition[] topicPartitions) { + TopicPartition[] topicPartitions) { this.consumerFactory = consumerFactory; this.topics = topics; this.topicPattern = topicPattern; @@ -203,7 +203,7 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener private class ListenerConsumer implements SchedulingAwareRunnable { - private final Log logger = LogFactory.getLog(this.getClass()); + private final Log logger = LogFactory.getLog(ListenerConsumer.class); private final CommitCallback callback = new CommitCallback(); @@ -221,7 +221,7 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener private final long recentOffset; - private final boolean autoCommit = consumerFactory.isAutoCommit(); + private final boolean autoCommit = KafkaMessageListenerContainer.this.consumerFactory.isAutoCommit(); private Thread consumerThread; @@ -230,35 +230,35 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener private volatile Collection assignedPartitions; public ListenerConsumer(MessageListener listener, AcknowledgingMessageListener ackListener, - ContainerOffsetResetStrategy resetStrategy, long recentOffset) { + ContainerOffsetResetStrategy resetStrategy, long recentOffset) { Assert.state(!(getAckMode().equals(AckMode.MANUAL) || getAckMode().equals(AckMode.MANUAL_IMMEDIATE)) || !this.autoCommit, "Consumer cannot be configured for auto commit for ackMode " + getAckMode()); - Consumer consumer = consumerFactory.createConsumer(); + Consumer consumer = KafkaMessageListenerContainer.this.consumerFactory.createConsumer(); ConsumerRebalanceListener rebalanceListener = new ConsumerRebalanceListener() { @Override public void onPartitionsRevoked(Collection partitions) { - logger.info("partitions revoked:" + partitions); + KafkaMessageListenerContainer.this.logger.info("partitions revoked:" + partitions); } @Override public void onPartitionsAssigned(Collection partitions) { - assignedPartitions = partitions; - logger.info("partitions assigned:" + partitions); + ListenerConsumer.this.assignedPartitions = partitions; + KafkaMessageListenerContainer.this.logger.info("partitions assigned:" + partitions); } }; - if (partitions == null) { - if (topicPattern != null) { - consumer.subscribe(topicPattern, rebalanceListener); + if (KafkaMessageListenerContainer.this.partitions == null) { + if (KafkaMessageListenerContainer.this.topicPattern != null) { + consumer.subscribe(KafkaMessageListenerContainer.this.topicPattern, rebalanceListener); } else { - consumer.subscribe(Arrays.asList(topics), rebalanceListener); + consumer.subscribe(Arrays.asList(KafkaMessageListenerContainer.this.topics), rebalanceListener); } } else { - List topicPartitions = Arrays.asList(partitions); + List topicPartitions = Arrays.asList(KafkaMessageListenerContainer.this.partitions); this.definedPartitions = topicPartitions; consumer.assign(topicPartitions); } @@ -280,26 +280,26 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener int count = 0; long last = System.currentTimeMillis(); long now; - if (isRunning() && definedPartitions != null) { + if (isRunning() && this.definedPartitions != null) { initPartitionsIfNeeded(); } final AckMode ackMode = getAckMode(); while (isRunning()) { try { - if (logger.isTraceEnabled()) { - logger.trace("Polling..."); + if (this.logger.isTraceEnabled()) { + this.logger.trace("Polling..."); } - ConsumerRecords records = consumer.poll(getPollTimeout()); + ConsumerRecords records = this.consumer.poll(getPollTimeout()); if (records != null) { count += records.count(); - if (logger.isDebugEnabled()) { - logger.debug("Received: " + records.count() + " records"); + if (this.logger.isDebugEnabled()) { + this.logger.debug("Received: " + records.count() + " records"); } Iterator> iterator = records.iterator(); while (iterator.hasNext()) { final ConsumerRecord record = iterator.next(); invokeListener(record); - if (!autoCommit && ackMode.equals(AckMode.RECORD)) { + if (!this.autoCommit && ackMode.equals(AckMode.RECORD)) { this.consumer.commitAsync( Collections.singletonMap(new TopicPartition(record.topic(), record.partition()), new OffsetAndMetadata(record.offset() + 1)), this.callback); @@ -338,35 +338,35 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener } } else { - if (logger.isDebugEnabled()) { - logger.debug("No records"); + if (this.logger.isDebugEnabled()) { + this.logger.debug("No records"); } } } catch (WakeupException e) { - ; + // No-op. Continue process } catch (Exception e) { if (getErrorHandler() != null) { getErrorHandler().handle(e, null); } else { - logger.error("Container exception", e); + this.logger.error("Container exception", e); } } } - if (offsets.size() > 0) { + if (this.offsets.size() > 0) { commitIfNecessary(); } try { this.consumer.unsubscribe(); } catch (WakeupException e) { - ; + // No-op. Continue process } this.consumer.close(); - if (logger.isInfoEnabled()) { - logger.info("Consumer stopped"); + if (this.logger.isInfoEnabled()) { + this.logger.info("Consumer stopped"); } } @@ -381,18 +381,18 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener updateManualOffset(record); } else if (getAckMode().equals(AckMode.MANUAL_IMMEDIATE)) { - if (Thread.currentThread().equals(consumerThread)) { + if (Thread.currentThread().equals(ListenerConsumer.this.consumerThread)) { Map commits = Collections.singletonMap( new TopicPartition(record.topic(), record.partition()), new OffsetAndMetadata(record.offset() + 1)); - if (logger.isDebugEnabled()) { - logger.debug("Committing: " + commits); + if (ListenerConsumer.this.logger.isDebugEnabled()) { + ListenerConsumer.this.logger.debug("Committing: " + commits); } - consumer.commitAsync(commits, callback); + ListenerConsumer.this.consumer.commitAsync(commits, ListenerConsumer.this.callback); } else { throw new IllegalStateException( - "With MANUAL_IMMEDIATE ack mode, acknowledget must be invoked on the " + "With MANUAL_IMMEDIATE ack mode, acknowledge() must be invoked on the " + "consumer thread"); } } @@ -405,7 +405,7 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener }); } else { - listener.onMessage(record); + this.listener.onMessage(record); } } catch (Exception e) { @@ -413,7 +413,7 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener getErrorHandler().handle(e, record); } else { - logger.error("Listener threw an exception and no error handler for " + record, e); + this.logger.error("Listener threw an exception and no error handler for " + record, e); } } } @@ -438,8 +438,8 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener for (TopicPartition topicPartition : this.definedPartitions) { long newOffset = this.consumer.position(topicPartition) - this.recentOffset; this.consumer.seek(topicPartition, newOffset); - if (logger.isDebugEnabled()) { - logger.debug("Reset " + topicPartition + " to offset " + newOffset); + if (this.logger.isDebugEnabled()) { + this.logger.debug("Reset " + topicPartition + " to offset " + newOffset); } } } @@ -483,8 +483,8 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener } } this.offsets.clear(); - if (logger.isDebugEnabled()) { - logger.debug("Committing: " + commits); + if (this.logger.isDebugEnabled()) { + this.logger.debug("Committing: " + commits); } if (!commits.isEmpty()) { this.consumer.commitAsync(commits, this.callback); @@ -494,7 +494,7 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener private static final class CommitCallback implements OffsetCommitCallback { - private final Log logger = LogFactory.getLog(OffsetCommitCallback.class); + private static final Log logger = LogFactory.getLog(OffsetCommitCallback.class); @Override public void onComplete(Map offsets, Exception exception) { diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/ListenerExecutionFailedException.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/ListenerExecutionFailedException.java index d2368a49..f39ba64d 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/ListenerExecutionFailedException.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/ListenerExecutionFailedException.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.listener; import org.springframework.kafka.core.KafkaException; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/LoggingErrorHandler.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/LoggingErrorHandler.java index 5cd133c0..4c78813e 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/LoggingErrorHandler.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/LoggingErrorHandler.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -14,7 +14,6 @@ * limitations under the License. */ - package org.springframework.kafka.listener; import org.apache.commons.logging.Log; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListener.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListener.java index 9576a55c..b5d5cdda 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListener.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListener.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -14,7 +14,6 @@ * limitations under the License. */ - package org.springframework.kafka.listener; import org.apache.kafka.clients.consumer.ConsumerRecord; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/MethodKafkaListenerEndpoint.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/MethodKafkaListenerEndpoint.java index f3fabc0a..45e47ca9 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/MethodKafkaListenerEndpoint.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/MethodKafkaListenerEndpoint.java @@ -79,7 +79,7 @@ public class MethodKafkaListenerEndpoint extends AbstractKafkaListenerEndp * @return the messageHandlerMethodFactory */ protected MessageHandlerMethodFactory getMessageHandlerMethodFactory() { - return messageHandlerMethodFactory; + return this.messageHandlerMethodFactory; } @Override diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/MultiMethodKafkaListenerEndpoint.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/MultiMethodKafkaListenerEndpoint.java index 27c11393..ff8bc3b6 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/MultiMethodKafkaListenerEndpoint.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/MultiMethodKafkaListenerEndpoint.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.listener; import java.lang.reflect.Method; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/AbstractAdaptableMessageListener.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/AbstractAdaptableMessageListener.java index 05de7d33..29fabfad 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/AbstractAdaptableMessageListener.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/AbstractAdaptableMessageListener.java @@ -71,7 +71,7 @@ public abstract class AbstractAdaptableMessageListener implements MessageL * @see #onMessage(ConsumerRecord) */ protected void handleListenerException(Throwable ex) { - logger.error("Listener execution failed", ex); + this.logger.error("Listener execution failed", ex); } /** diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/DelegatingInvocableHandler.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/DelegatingInvocableHandler.java index e3c7cc5e..8a7071e0 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/DelegatingInvocableHandler.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/DelegatingInvocableHandler.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.listener.adapter; import java.lang.annotation.Annotation; @@ -60,7 +61,7 @@ public class DelegatingInvocableHandler { * @return the bean */ public Object getBean() { - return bean; + return this.bean; } /** @@ -88,7 +89,7 @@ public class DelegatingInvocableHandler { if (handler == null) { throw new KafkaException("No method found for " + payloadClass); } - this.cachedHandlers.putIfAbsent(payloadClass, handler);//NOSONAR + this.cachedHandlers.putIfAbsent(payloadClass, handler); //NOSONAR } return handler; } @@ -141,7 +142,7 @@ public class DelegatingInvocableHandler { */ public String getMethodNameFor(Object payload) { InvocableHandlerMethod handlerForPayload = getHandlerForPayload(payload.getClass()); - return handlerForPayload == null ? "no match" : handlerForPayload.getMethod().toGenericString();//NOSONAR + return handlerForPayload == null ? "no match" : handlerForPayload.getMethod().toGenericString(); //NOSONAR } } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/HandlerAdapter.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/HandlerAdapter.java index 84b79b5b..52d0f2fd 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/HandlerAdapter.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/HandlerAdapter.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.listener.adapter; import org.springframework.messaging.Message; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessageConverter.java b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessageConverter.java index 48cd3d93..ac650939 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessageConverter.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessageConverter.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.support.converter; import org.apache.kafka.clients.consumer.ConsumerRecord; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessagingMessageConverter.java b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessagingMessageConverter.java index 11dcf157..7b09580f 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessagingMessageConverter.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessagingMessageConverter.java @@ -56,7 +56,7 @@ public class MessagingMessageConverter implements MessageConverter { @Override public Message toMessage(ConsumerRecord record, Acknowledgment acknowledgment) { - KafkaMessageHeaders kafkaMessageHeaders = new KafkaMessageHeaders(generateMessageId, generateTimestamp); + KafkaMessageHeaders kafkaMessageHeaders = new KafkaMessageHeaders(this.generateMessageId, this.generateTimestamp); Map rawHeaders = kafkaMessageHeaders.getRawHeaders(); rawHeaders.put(KafkaHeaders.MESSAGE_KEY, record.key()); diff --git a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java index 137a2699..74fff6b7 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.annotation; import static org.junit.Assert.assertEquals; @@ -116,7 +117,7 @@ public class EnableKafkaIntegrationTests { @Bean public KafkaListenerContainerFactory> - kafkaListenerContainerFactory() { + kafkaListenerContainerFactory() { SimpleKafkaListenerContainerFactory factory = new SimpleKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); return factory; @@ -124,7 +125,7 @@ public class EnableKafkaIntegrationTests { @Bean public KafkaListenerContainerFactory> - kafkaManualAckListenerContainerFactory() { + kafkaManualAckListenerContainerFactory() { SimpleKafkaListenerContainerFactory factory = new SimpleKafkaListenerContainerFactory<>(); factory.setConsumerFactory(manualConsumerFactory()); factory.setAckMode(AckMode.MANUAL_IMMEDIATE); @@ -204,24 +205,24 @@ public class EnableKafkaIntegrationTests { private volatile Acknowledgment ack; - @KafkaListener(id="foo", topics = "annotated1") + @KafkaListener(id = "foo", topics = "annotated1") public void listen1(String foo) { this.latch1.countDown(); } - @KafkaListener(id="bar", topicPattern = "annotated2") + @KafkaListener(id = "bar", topicPattern = "annotated2") public void listen2(@Payload String foo, @Header(KafkaHeaders.PARTITION_ID) int partitionHeader) { this.partition = partitionHeader; this.latch2.countDown(); } - @KafkaListener(id="baz", topicPartitions = @TopicPartition(topic = "annotated3", partition="0")) + @KafkaListener(id = "baz", topicPartitions = @TopicPartition(topic = "annotated3", partition = "0")) public void listen3(ConsumerRecord record) { this.record = record; this.latch3.countDown(); } - @KafkaListener(id="qux", topics = "annotated4", containerFactory = "kafkaManualAckListenerContainerFactory") + @KafkaListener(id = "qux", topics = "annotated4", containerFactory = "kafkaManualAckListenerContainerFactory") public void listen4(@Payload String foo, Acknowledgment ack) { this.ack = ack; this.ack.acknowledge(); diff --git a/spring-kafka/src/test/java/org/springframework/kafka/core/BrokerAddress.java b/spring-kafka/src/test/java/org/springframework/kafka/core/BrokerAddress.java index 72e7a571..0df06b74 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/core/BrokerAddress.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/core/BrokerAddress.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -14,7 +14,6 @@ * limitations under the License. */ - package org.springframework.kafka.core; import org.springframework.util.Assert; diff --git a/src/checkstyle/checkstyle-header.txt b/src/checkstyle/checkstyle-header.txt new file mode 100644 index 00000000..f470e973 --- /dev/null +++ b/src/checkstyle/checkstyle-header.txt @@ -0,0 +1,17 @@ +^\Q/*\E$ +^\Q * Copyright \E20\d\d(\-20\d\d)?\Q the original author or authors.\E$ +^\Q *\E$ +^\Q * Licensed under the Apache License, Version 2.0 (the "License");\E$ +^\Q * you may not use this file except in compliance with the License.\E$ +^\Q * You may obtain a copy of the License at\E$ +^\Q *\E$ +^\Q * http://www.apache.org/licenses/LICENSE-2.0\E$ +^\Q *\E$ +^\Q * Unless required by applicable law or agreed to in writing, software\E$ +^\Q * distributed under the License is distributed on an "AS IS" BASIS,\E$ +^\Q * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.\E$ +^\Q * See the License for the specific language governing permissions and\E$ +^\Q * limitations under the License.\E$ +^\Q */\E$ +^$ +^.*$ diff --git a/src/checkstyle/checkstyle-suppressions.xml b/src/checkstyle/checkstyle-suppressions.xml new file mode 100644 index 00000000..dcc709e4 --- /dev/null +++ b/src/checkstyle/checkstyle-suppressions.xml @@ -0,0 +1,8 @@ + + + + + + diff --git a/src/checkstyle/checkstyle.xml b/src/checkstyle/checkstyle.xml new file mode 100644 index 00000000..3453f77f --- /dev/null +++ b/src/checkstyle/checkstyle.xml @@ -0,0 +1,169 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/src/reference/asciidoc/quick-tour.adoc b/src/reference/asciidoc/quick-tour.adoc index 8611a426..17246dd9 100644 --- a/src/reference/asciidoc/quick-tour.adoc +++ b/src/reference/asciidoc/quick-tour.adoc @@ -72,7 +72,8 @@ public void testAutoCommit() throws Exception { private KafkaMessageListenerContainer createContainer() { Map props = consumerProps(); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory(props); - KafkaMessageListenerContainer container = new KafkaMessageListenerContainer<>(cf, topic1); + KafkaMessageListenerContainer container = + new KafkaMessageListenerContainer<>(cf, topic1); return container; } @@ -137,7 +138,8 @@ public class Config { @Bean public KafkaListenerContainerFactory> kafkaListenerContainerFactory() { - SimpleKafkaListenerContainerFactory factory = new SimpleKafkaListenerContainerFactory<>(); + SimpleKafkaListenerContainerFactory factory = + new SimpleKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); return factory; } @@ -151,13 +153,7 @@ public class Config { public Map consumerConfigs() { Map props = new HashMap<>(); props.put("bootstrap.servers", embeddedKafka.getBrokersAsString()); - props.put("bootstrap.servers", "localhost:9092"); - props.put("group.id", "myGroup"); - props.put("enable.auto.commit", true); - props.put("auto.commit.interval.ms", "100"); - props.put("session.timeout.ms", "15000"); - props.put("key.deserializer", "org.apache.kafka.common.serialization.IntegerDeserializer"); - props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); + ...... return props; } @@ -175,13 +171,7 @@ public class Config { public Map producerConfigs() { Map props = new HashMap<>(); props.put("bootstrap.servers", embeddedKafka.getBrokersAsString()); - props.put("bootstrap.servers", "localhost:9092"); - props.put("retries", 0); - props.put("batch.size", 16384); - props.put("linger.ms", 1); - props.put("buffer.memory", 33554432); - props.put("key.serializer", "org.apache.kafka.common.serialization.IntegerSerializer"); - props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); + ...... return props; } @@ -196,7 +186,7 @@ public class Listener { private final CountDownLatch latch1 = new CountDownLatch(1); - @KafkaListener(id="foo", topics = "annotated1") + @KafkaListener(id = "foo", topics = "annotated1") public void listen1(String foo) { this.latch1.countDown(); }