Skip to content

Commit 77e4112

Browse files
[Pull-based Ingestion] Fix snappy related Kafka errors and relax thread leak checks (#17240)
* Fix snappy related Kafka errors and relax thread leak checks Signed-off-by: Varun Bharadwaj <[email protected]> * update comment Signed-off-by: Varun Bharadwaj <[email protected]> * Skip thread leak checks in kafka IT Signed-off-by: Varun Bharadwaj <[email protected]> * exclude testcontainer watchdog thread from thread leak checks Signed-off-by: Varun Bharadwaj <[email protected]> --------- Signed-off-by: Varun Bharadwaj <[email protected]>
1 parent 8089b61 commit 77e4112

File tree

7 files changed

+288
-27
lines changed

7 files changed

+288
-27
lines changed

plugins/ingestion-kafka/build.gradle

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,13 +21,17 @@ versions << [
2121
'docker': '3.3.6',
2222
'testcontainers': '1.19.7',
2323
'ducttape': '1.0.8',
24+
'snappy': '1.1.10.7',
2425
]
2526

2627
dependencies {
2728
// kafka
2829
api "org.slf4j:slf4j-api:${versions.slf4j}"
2930
api "org.apache.kafka:kafka-clients:${versions.kafka}"
3031

32+
// dependencies of kafka-clients
33+
runtimeOnly "org.xerial.snappy:snappy-java:${versions.snappy}"
34+
3135
// test
3236
testImplementation "com.github.docker-java:docker-java-api:${versions.docker}"
3337
testImplementation "com.github.docker-java:docker-java-transport:${versions.docker}"
@@ -107,6 +111,8 @@ thirdPartyAudit {
107111
'org.jose4j.jwt.consumer.JwtContext',
108112
'org.jose4j.jwx.Headers',
109113
'org.jose4j.keys.resolvers.VerificationKeyResolver',
114+
'org.osgi.framework.BundleActivator',
115+
'org.osgi.framework.BundleContext',
110116
)
111117
ignoreViolations(
112118
'org.apache.kafka.shaded.com.google.protobuf.MessageSchema',
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
3049f95640f4625a945cfab85715f603fa4c8f80
Lines changed: 202 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,202 @@
1+
2+
Apache License
3+
Version 2.0, January 2004
4+
http://www.apache.org/licenses/
5+
6+
TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION
7+
8+
1. Definitions.
9+
10+
"License" shall mean the terms and conditions for use, reproduction,
11+
and distribution as defined by Sections 1 through 9 of this document.
12+
13+
"Licensor" shall mean the copyright owner or entity authorized by
14+
the copyright owner that is granting the License.
15+
16+
"Legal Entity" shall mean the union of the acting entity and all
17+
other entities that control, are controlled by, or are under common
18+
control with that entity. For the purposes of this definition,
19+
"control" means (i) the power, direct or indirect, to cause the
20+
direction or management of such entity, whether by contract or
21+
otherwise, or (ii) ownership of fifty percent (50%) or more of the
22+
outstanding shares, or (iii) beneficial ownership of such entity.
23+
24+
"You" (or "Your") shall mean an individual or Legal Entity
25+
exercising permissions granted by this License.
26+
27+
"Source" form shall mean the preferred form for making modifications,
28+
including but not limited to software source code, documentation
29+
source, and configuration files.
30+
31+
"Object" form shall mean any form resulting from mechanical
32+
transformation or translation of a Source form, including but
33+
not limited to compiled object code, generated documentation,
34+
and conversions to other media types.
35+
36+
"Work" shall mean the work of authorship, whether in Source or
37+
Object form, made available under the License, as indicated by a
38+
copyright notice that is included in or attached to the work
39+
(an example is provided in the Appendix below).
40+
41+
"Derivative Works" shall mean any work, whether in Source or Object
42+
form, that is based on (or derived from) the Work and for which the
43+
editorial revisions, annotations, elaborations, or other modifications
44+
represent, as a whole, an original work of authorship. For the purposes
45+
of this License, Derivative Works shall not include works that remain
46+
separable from, or merely link (or bind by name) to the interfaces of,
47+
the Work and Derivative Works thereof.
48+
49+
"Contribution" shall mean any work of authorship, including
50+
the original version of the Work and any modifications or additions
51+
to that Work or Derivative Works thereof, that is intentionally
52+
submitted to Licensor for inclusion in the Work by the copyright owner
53+
or by an individual or Legal Entity authorized to submit on behalf of
54+
the copyright owner. For the purposes of this definition, "submitted"
55+
means any form of electronic, verbal, or written communication sent
56+
to the Licensor or its representatives, including but not limited to
57+
communication on electronic mailing lists, source code control systems,
58+
and issue tracking systems that are managed by, or on behalf of, the
59+
Licensor for the purpose of discussing and improving the Work, but
60+
excluding communication that is conspicuously marked or otherwise
61+
designated in writing by the copyright owner as "Not a Contribution."
62+
63+
"Contributor" shall mean Licensor and any individual or Legal Entity
64+
on behalf of whom a Contribution has been received by Licensor and
65+
subsequently incorporated within the Work.
66+
67+
2. Grant of Copyright License. Subject to the terms and conditions of
68+
this License, each Contributor hereby grants to You a perpetual,
69+
worldwide, non-exclusive, no-charge, royalty-free, irrevocable
70+
copyright license to reproduce, prepare Derivative Works of,
71+
publicly display, publicly perform, sublicense, and distribute the
72+
Work and such Derivative Works in Source or Object form.
73+
74+
3. Grant of Patent License. Subject to the terms and conditions of
75+
this License, each Contributor hereby grants to You a perpetual,
76+
worldwide, non-exclusive, no-charge, royalty-free, irrevocable
77+
(except as stated in this section) patent license to make, have made,
78+
use, offer to sell, sell, import, and otherwise transfer the Work,
79+
where such license applies only to those patent claims licensable
80+
by such Contributor that are necessarily infringed by their
81+
Contribution(s) alone or by combination of their Contribution(s)
82+
with the Work to which such Contribution(s) was submitted. If You
83+
institute patent litigation against any entity (including a
84+
cross-claim or counterclaim in a lawsuit) alleging that the Work
85+
or a Contribution incorporated within the Work constitutes direct
86+
or contributory patent infringement, then any patent licenses
87+
granted to You under this License for that Work shall terminate
88+
as of the date such litigation is filed.
89+
90+
4. Redistribution. You may reproduce and distribute copies of the
91+
Work or Derivative Works thereof in any medium, with or without
92+
modifications, and in Source or Object form, provided that You
93+
meet the following conditions:
94+
95+
(a) You must give any other recipients of the Work or
96+
Derivative Works a copy of this License; and
97+
98+
(b) You must cause any modified files to carry prominent notices
99+
stating that You changed the files; and
100+
101+
(c) You must retain, in the Source form of any Derivative Works
102+
that You distribute, all copyright, patent, trademark, and
103+
attribution notices from the Source form of the Work,
104+
excluding those notices that do not pertain to any part of
105+
the Derivative Works; and
106+
107+
(d) If the Work includes a "NOTICE" text file as part of its
108+
distribution, then any Derivative Works that You distribute must
109+
include a readable copy of the attribution notices contained
110+
within such NOTICE file, excluding those notices that do not
111+
pertain to any part of the Derivative Works, in at least one
112+
of the following places: within a NOTICE text file distributed
113+
as part of the Derivative Works; within the Source form or
114+
documentation, if provided along with the Derivative Works; or,
115+
within a display generated by the Derivative Works, if and
116+
wherever such third-party notices normally appear. The contents
117+
of the NOTICE file are for informational purposes only and
118+
do not modify the License. You may add Your own attribution
119+
notices within Derivative Works that You distribute, alongside
120+
or as an addendum to the NOTICE text from the Work, provided
121+
that such additional attribution notices cannot be construed
122+
as modifying the License.
123+
124+
You may add Your own copyright statement to Your modifications and
125+
may provide additional or different license terms and conditions
126+
for use, reproduction, or distribution of Your modifications, or
127+
for any such Derivative Works as a whole, provided Your use,
128+
reproduction, and distribution of the Work otherwise complies with
129+
the conditions stated in this License.
130+
131+
5. Submission of Contributions. Unless You explicitly state otherwise,
132+
any Contribution intentionally submitted for inclusion in the Work
133+
by You to the Licensor shall be under the terms and conditions of
134+
this License, without any additional terms or conditions.
135+
Notwithstanding the above, nothing herein shall supersede or modify
136+
the terms of any separate license agreement you may have executed
137+
with Licensor regarding such Contributions.
138+
139+
6. Trademarks. This License does not grant permission to use the trade
140+
names, trademarks, service marks, or product names of the Licensor,
141+
except as required for reasonable and customary use in describing the
142+
origin of the Work and reproducing the content of the NOTICE file.
143+
144+
7. Disclaimer of Warranty. Unless required by applicable law or
145+
agreed to in writing, Licensor provides the Work (and each
146+
Contributor provides its Contributions) on an "AS IS" BASIS,
147+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
148+
implied, including, without limitation, any warranties or conditions
149+
of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A
150+
PARTICULAR PURPOSE. You are solely responsible for determining the
151+
appropriateness of using or redistributing the Work and assume any
152+
risks associated with Your exercise of permissions under this License.
153+
154+
8. Limitation of Liability. In no event and under no legal theory,
155+
whether in tort (including negligence), contract, or otherwise,
156+
unless required by applicable law (such as deliberate and grossly
157+
negligent acts) or agreed to in writing, shall any Contributor be
158+
liable to You for damages, including any direct, indirect, special,
159+
incidental, or consequential damages of any character arising as a
160+
result of this License or out of the use or inability to use the
161+
Work (including but not limited to damages for loss of goodwill,
162+
work stoppage, computer failure or malfunction, or any and all
163+
other commercial damages or losses), even if such Contributor
164+
has been advised of the possibility of such damages.
165+
166+
9. Accepting Warranty or Additional Liability. While redistributing
167+
the Work or Derivative Works thereof, You may choose to offer,
168+
and charge a fee for, acceptance of support, warranty, indemnity,
169+
or other liability obligations and/or rights consistent with this
170+
License. However, in accepting such obligations, You may act only
171+
on Your own behalf and on Your sole responsibility, not on behalf
172+
of any other Contributor, and only if You agree to indemnify,
173+
defend, and hold each Contributor harmless for any liability
174+
incurred by, or claims asserted against, such Contributor by reason
175+
of your accepting any such warranty or additional liability.
176+
177+
END OF TERMS AND CONDITIONS
178+
179+
APPENDIX: How to apply the Apache License to your work.
180+
181+
To apply the Apache License to your work, attach the following
182+
boilerplate notice, with the fields enclosed by brackets "[]"
183+
replaced with your own identifying information. (Don't include
184+
the brackets!) The text should be enclosed in the appropriate
185+
comment syntax for the file format. We also recommend that a
186+
file or class name and description of purpose be included on the
187+
same "printed page" as the copyright notice for easier
188+
identification within third-party archives.
189+
190+
Copyright [yyyy] [name of copyright owner]
191+
192+
Licensed under the Apache License, Version 2.0 (the "License");
193+
you may not use this file except in compliance with the License.
194+
You may obtain a copy of the License at
195+
196+
http://www.apache.org/licenses/LICENSE-2.0
197+
198+
Unless required by applicable law or agreed to in writing, software
199+
distributed under the License is distributed on an "AS IS" BASIS,
200+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
201+
See the License for the specific language governing permissions and
202+
limitations under the License.
Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,22 @@
1+
This product includes software developed by Google
2+
Snappy: http://code.google.com/p/snappy/ (New BSD License)
3+
4+
This product includes software developed by Apache
5+
PureJavaCrc32C from apache-hadoop-common http://hadoop.apache.org/
6+
(Apache 2.0 license)
7+
8+
This library contains statically linked libstdc++. This inclusion is allowed by
9+
"GCC Runtime Library Exception"
10+
http://gcc.gnu.org/onlinedocs/libstdc++/manual/license.html
11+
12+
== Contributors ==
13+
* Tatu Saloranta
14+
* Providing benchmark suite
15+
* Alec Wysoker
16+
* Performance and memory usage improvement
17+
18+
Third-Party Notices and Licenses:
19+
20+
- Hadoop: Apache Hadoop is used as a dependency
21+
License: Apache License 2.0
22+
Source/Reference: https://github.com/apache/hadoop/blob/trunk/NOTICE.txt

plugins/ingestion-kafka/src/internalClusterTest/java/org/opensearch/plugin/kafka/IngestFromKafkaIT.java

Lines changed: 30 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@
88

99
package org.opensearch.plugin.kafka;
1010

11-
import com.carrotsearch.randomizedtesting.annotations.ThreadLeakLingering;
11+
import com.carrotsearch.randomizedtesting.annotations.ThreadLeakFilters;
1212

1313
import org.apache.kafka.clients.producer.KafkaProducer;
1414
import org.apache.kafka.clients.producer.Producer;
@@ -45,7 +45,7 @@
4545
/**
4646
* Integration test for Kafka ingestion
4747
*/
48-
@ThreadLeakLingering(linger = 15000) // wait for container pull thread to die
48+
@ThreadLeakFilters(filters = TestContainerWatchdogThreadLeakFilter.class)
4949
public class IngestFromKafkaIT extends OpenSearchIntegTestCase {
5050
static final String topicName = "test";
5151

@@ -75,29 +75,31 @@ public void testPluginsAreInstalled() {
7575
}
7676

7777
public void testKafkaIngestion() {
78-
setupKafka();
79-
// create an index with ingestion source from kafka
80-
createIndex(
81-
"test",
82-
Settings.builder()
83-
.put(IndexMetadata.SETTING_NUMBER_OF_SHARDS, 1)
84-
.put(IndexMetadata.SETTING_NUMBER_OF_REPLICAS, 0)
85-
.put("ingestion_source.type", "kafka")
86-
.put("ingestion_source.pointer.init.reset", "earliest")
87-
.put("ingestion_source.param.topic", "test")
88-
.put("ingestion_source.param.bootstrap_servers", kafka.getBootstrapServers())
89-
.build(),
90-
"{\"properties\":{\"name\":{\"type\": \"text\"},\"age\":{\"type\": \"integer\"}}}}"
91-
);
92-
93-
RangeQueryBuilder query = new RangeQueryBuilder("age").gte(21);
94-
await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> {
95-
refresh("test");
96-
SearchResponse response = client().prepareSearch("test").setQuery(query).get();
97-
assertThat(response.getHits().getTotalHits().value(), is(1L));
98-
});
99-
100-
stopKafka();
78+
try {
79+
setupKafka();
80+
// create an index with ingestion source from kafka
81+
createIndex(
82+
"test",
83+
Settings.builder()
84+
.put(IndexMetadata.SETTING_NUMBER_OF_SHARDS, 1)
85+
.put(IndexMetadata.SETTING_NUMBER_OF_REPLICAS, 0)
86+
.put("ingestion_source.type", "kafka")
87+
.put("ingestion_source.pointer.init.reset", "earliest")
88+
.put("ingestion_source.param.topic", "test")
89+
.put("ingestion_source.param.bootstrap_servers", kafka.getBootstrapServers())
90+
.build(),
91+
"{\"properties\":{\"name\":{\"type\": \"text\"},\"age\":{\"type\": \"integer\"}}}}"
92+
);
93+
94+
RangeQueryBuilder query = new RangeQueryBuilder("age").gte(21);
95+
await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> {
96+
refresh("test");
97+
SearchResponse response = client().prepareSearch("test").setQuery(query).get();
98+
assertThat(response.getHits().getTotalHits().value(), is(1L));
99+
});
100+
} finally {
101+
stopKafka();
102+
}
101103
}
102104

103105
private void setupKafka() {
@@ -109,7 +111,9 @@ private void setupKafka() {
109111
}
110112

111113
private void stopKafka() {
112-
kafka.stop();
114+
if (kafka != null) {
115+
kafka.stop();
116+
}
113117
}
114118

115119
private void prepareKafkaData() {
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,23 @@
1+
/*
2+
* SPDX-License-Identifier: Apache-2.0
3+
*
4+
* The OpenSearch Contributors require contributions made to
5+
* this file be licensed under the Apache-2.0 license or a
6+
* compatible open source license.
7+
*/
8+
9+
package org.opensearch.plugin.kafka;
10+
11+
import com.carrotsearch.randomizedtesting.ThreadFilter;
12+
13+
/**
14+
* The {@link org.testcontainers.images.TimeLimitedLoggedPullImageResultCallback} instance used by test containers,
15+
* for example {@link org.testcontainers.containers.KafkaContainer} creates a watcher daemon thread which is never
16+
* stopped. This filter excludes that thread from the thread leak detection logic.
17+
*/
18+
public final class TestContainerWatchdogThreadLeakFilter implements ThreadFilter {
19+
@Override
20+
public boolean reject(Thread t) {
21+
return t.getName().startsWith("testcontainers-pull-watchdog-");
22+
}
23+
}

plugins/ingestion-kafka/src/main/plugin-metadata/plugin-security.policy

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,5 +17,8 @@ grant {
1717
// Allow host/ip name service lookups
1818
permission java.net.SocketPermission "*", "connect";
1919
permission java.net.SocketPermission "*", "resolve";
20-
};
2120

21+
// Needed for Kafka consumer to load native snappy library
22+
permission java.lang.RuntimePermission "loadLibrary.*";
23+
permission java.io.FilePermission "${java.io.tmpdir}${/}snappy-*", "read,write";
24+
};

0 commit comments

Comments
 (0)