-
Provides an efficient custom Teragrep datasource for Spark
-
Combines results from Kafka and from the S3
-
Configurable
-
Metrics
See the official documentation on docs.teragrep.com.
-
Designed on Java 8 and confirmed to also work on Java 11.
-
Core spark and aws features are marked as provided to keep project lightweight and ready to deploy to a spark cluster.
We start by adding the project to our version control, here is an example with Maven:
<dependency>
<groupId>com.teragrep</groupId>
<artifactId>pth_06</artifactId>
<version>4.8.2</version> <!-- use latest version -->
<scope>compile</scope>
</dependency>We can then access the custom Spark datasource, here is a Java example:
import com.teragrep.pth_06.TeragrepDatasource;
Dataset<Row> df = spark.readStream()
.format("com.teragrep.pth_06.TeragrepDatasource")
.option("archive.enabled", "true")
.option("kafka.enabled", "true")
.option("TeragrepAuditQuery", "<index value=\"example\" operation=\"EQUALS\"");How do we query the wanted results from the archive? The datasource can accept a query XML that will be used to fetch the wanted results.
The value tags should have three values:
-
tag name
-
value
-
comparison operation
Logical elements (<AND>, <OR>) are used to logically combine value tag results, you can have multiple of either tags.
Here is an example of an XML where we query for results that are in the index example and after the epoch 1643207821.
<AND>
<index value="example" operation="EQUALS"/>
<earliest operation="GE" value="1643207821"/>
</AND>You can use <AND> and <OR> to create very specific queries.
Currently supported XML Elements and their supported operations.
| Element | Purpose | Operations |
|---|---|---|
<index> |
filter by index |
EQUALS, NOT_EQUALS |
<earliest> or <index_earliest> |
start time of the query |
EQUALS |
<latest> or <index_latest> |
end time of the query |
EQUALS |
<sourcetype> |
filter by sourcetype |
EQUALS, NOT_EQUALS |
<host> |
filter by host |
EQUALS, NOT_EQUALS |
<indexstatement> |
keyword filter, this is required only when using bloomfilter accelerated searches |
EQUALS |
Configures the common options
| Option | Purpose | Type |
|---|---|---|
archive.enabled |
Enables long-term archive results. |
Boolean |
kafka.enabled |
Enables near real-time kafka results. |
Boolean |
S3endPoint |
S3 end point. |
Required String |
S3identity |
S3 identity. |
Required String |
S3credential |
S3 credential. |
Required String |
DBusername |
Database username. |
Required String |
DBpassword |
Database password. |
Required String |
DBurl |
Database url. |
Required String |
bloom.enabled |
Enable Bloom filter pattern acceleration, requires bloom filter creation using PTH_10. |
Boolean |
bloom.withoutFilters |
Enable option to filter results based if they have a filter with a certain pattern created. Admin command. |
Boolean |
bloom.withoutFiltersPattern |
The pattern that is used to check if a filter with that pattern has been created for the result. Admin command. |
String |
DBbloomdbname |
Bloom database name. |
String |
DBjournaldbname |
Journal database name. |
String |
DBstreamdbname |
Stream database name. |
String |
skipNonRFC5424Files |
Enables skipping of results that are not in the RFC5424 standard. |
Boolean |
epochMigrationMode |
Enables epoch migration mode, used by admin to migrate missing epoch values to the metadata. Admin command. |
Boolean |
archive.includeBeforeEpoch |
Option to precede the query earliest value by a amount of time, i.e. including previous timezones or wrong epoch values. Admin command. |
Long |
Configures the auditing information for the query, provides a way to log the reason why the query is made and by whom.
| Option | Purpose | Type |
|---|---|---|
TeragrepAuditQuery |
Query XML string used to determine the results. |
String |
TeragrepAuditReason |
Reason for making the query, i.e. why did I access this dataset. |
String |
TeragrepAuditUser |
User ID who is executing the query. |
String |
TeragrepAuditPluginClassName |
Implementation of RAD_01 to provide further audit information for the query. Including information about who accessed which records and what data did they contain. This is to satisfy the full audit trail requirement. |
String |
Configures the batch size settings, used for fine-tuning performance.
| Option | Purpose | Type | Recommended Default |
|---|---|---|---|
num_partitions |
Number of partitions in Spark. |
Integer |
24 |
quantumLength |
Base weight factor used to calculate maximum data weight per task. |
Integer |
15 |
batch.size.fileCompressionRatio |
Average Compression rate of the S3 files after the archive process. |
Float |
15.5 |
batch.size.processingSpeed |
Processing speed. |
Float |
136.5 |
batch.size.totalObjectCountLimit |
Limits the total maximum objects for batch. |
Long |
1000 |
Configures Kafka, you can also add common Kafka options like kafka.bootstrap.servers.
| Option | Purpose | Type |
|---|---|---|
kafka.continuousProcessing |
Enables continuous processing mode which means that the query will not self-terminate. |
Boolean |
skipNonRFC5424Files |
Enables skipping Kafka results that are not in the RFC5424 standard. |
Boolean |
Configures logging.
| Option | Purpose | Type |
|---|---|---|
logging.debug.enabled |
Enables debug logging. |
Boolean |
You can involve yourself with our project by opening an issue or submitting a pull request.
Contribution requirements:
-
All changes must be accompanied by a new or changed test. If you think testing is not required in your pull request, include a sufficient explanation as why you think so.
-
Security checks must pass
-
Pull requests must align with the principles and values of extreme programming.
-
Pull requests must follow the principles of Object Thinking and Elegant Objects (EO).
Read more in our Contributing Guideline.
Contributors must sign Teragrep Contributor License Agreement before a pull request is accepted to organization’s repositories.
You need to submit the CLA only once. After submitting the CLA you can contribute to all Teragrep’s repositories.