-
Notifications
You must be signed in to change notification settings - Fork 330
/
LzopCodec.java
183 lines (163 loc) · 6.82 KB
/
LzopCodec.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
/*
* This file is part of Hadoop-Gpl-Compression.
*
* Hadoop-Gpl-Compression is free software: you can redistribute it
* and/or modify it under the terms of the GNU General Public License
* as published by the Free Software Foundation, either version 3 of
* the License, or (at your option) any later version.
*
* Hadoop-Gpl-Compression is distributed in the hope that it will be
* useful, but WITHOUT ANY WARRANTY; without even the implied warranty
* of MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with Hadoop-Gpl-Compression. If not, see
* <https://www.gnu.org/licenses/>.
*/
package com.hadoop.compression.lzo;
import java.io.DataOutputStream;
import java.io.FilterInputStream;
import java.io.FilterOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.io.compress.CodecPool;
import org.apache.hadoop.io.compress.CompressionCodec;
import org.apache.hadoop.io.compress.CompressionInputStream;
import org.apache.hadoop.io.compress.CompressionOutputStream;
import org.apache.hadoop.io.compress.Compressor;
import org.apache.hadoop.io.compress.Decompressor;
/**
* A {@link CompressionCodec} for a streaming
* <b>lzo</b> compression/decompression pair compatible with lzop.
* https://www.lzop.org/
*/
public class LzopCodec extends LzoCodec {
/** 9 bytes at the top of every lzo file */
public static final byte[] LZO_MAGIC = new byte[] {
-119, 'L', 'Z', 'O', 0, '\r', '\n', '\032', '\n' };
/** Version of lzop this emulates */
public static final int LZOP_VERSION = 0x1010;
/** Latest verion of lzop this should be compatible with */
public static final int LZOP_COMPAT_VERSION = 0x0940;
public static final String DEFAULT_LZO_EXTENSION = ".lzo";
/**
* CodecPool.getCompressor() that takes conf is supported only in CDH3.
* The change is yet to make it to Apache Hadoop. Fall back to old
* getCompressor() if the new interface is not present.
*/
private static boolean codecPoolSupportsConf = false;
static {
try {
codecPoolSupportsConf =
null != CodecPool.class.getMethod("getCompressor",
CompressionCodec.class,
Configuration.class);
} catch (Exception e) {
}
}
@Override
public CompressionOutputStream createOutputStream(OutputStream out) throws IOException {
//get a compressor which will be returned to the pool when the output stream
//is closed.
Compressor compressor = getCompressor();
OutputStream wrapped = new WrappedOutputStream(out, compressor);
return createOutputStream(wrapped, compressor);
}
public CompressionOutputStream createIndexedOutputStream(OutputStream out,
DataOutputStream indexOut)
throws IOException {
//get a compressor which will be returned to the pool when the output stream
//is closed.
Compressor compressor = getCompressor();
OutputStream wrapped = new WrappedOutputStream(out, compressor);
return createIndexedOutputStream(wrapped, indexOut, compressor);
}
@Override
public CompressionOutputStream createOutputStream(OutputStream out,
Compressor compressor) throws IOException {
return createIndexedOutputStream(out, null, compressor);
}
public CompressionOutputStream createIndexedOutputStream(OutputStream out,
DataOutputStream indexOut, Compressor compressor) throws IOException {
if (!isNativeLzoLoaded(getConf())) {
throw new RuntimeException("native-lzo library not available");
}
LzoCompressor.CompressionStrategy strategy = LzoCompressor.CompressionStrategy.valueOf(
getConf().get(LZO_COMPRESSOR_KEY, LzoCompressor.CompressionStrategy.LZO1X_1.name()));
int bufferSize = getConf().getInt(LZO_BUFFER_SIZE_KEY, DEFAULT_LZO_BUFFER_SIZE);
return new LzopOutputStream(out, indexOut, compressor, bufferSize, strategy);
}
@Override
public CompressionInputStream createInputStream(InputStream in,
Decompressor decompressor) throws IOException {
// Ensure native-lzo library is loaded & initialized
if (!isNativeLzoLoaded(getConf())) {
throw new RuntimeException("native-lzo library not available");
}
return new LzopInputStream(in, decompressor,
getConf().getInt(LZO_BUFFER_SIZE_KEY, DEFAULT_LZO_BUFFER_SIZE));
}
// Previous versions of the API accidentally added/removed compressor/decompressors from the pool
// when they shouldn't. This classs is kind of a hack to maintain existing behavior,
// while still allowing proper resource management from outside
private static class WrappedOutputStream extends FilterOutputStream {
private Compressor compressor;
public WrappedOutputStream(OutputStream outputStream, Compressor compressor) {
super(outputStream);
this.compressor = compressor;
}
@Override
public void write(byte[] b, int off, int len) throws IOException {
out.write(b, off, len);
}
@Override
public void close() throws IOException {
CodecPool.returnCompressor(compressor);
super.close();
}
}
@Override
public CompressionInputStream createInputStream(final InputStream in) throws IOException {
final Decompressor decompressor = CodecPool.getDecompressor(this);
// maintain backwards compatibility re: returning the decompressor to the CodecPool
InputStream inputStream = new FilterInputStream(in) {
@Override
public void close() throws IOException {
CodecPool.returnDecompressor(decompressor);
super.close();
}
};
return createInputStream(inputStream, decompressor);
}
@Override
public Class<? extends Decompressor> getDecompressorType() {
// Ensure native-lzo library is loaded & initialized
if (!isNativeLzoLoaded(getConf())) {
throw new RuntimeException("native-lzo library not available");
}
return LzopDecompressor.class;
}
@Override
public Decompressor createDecompressor() {
if (!isNativeLzoLoaded(getConf())) {
throw new RuntimeException("native-lzo library not available");
}
return new LzopDecompressor(getConf().getInt(LZO_BUFFER_SIZE_KEY, DEFAULT_LZO_BUFFER_SIZE));
}
private Compressor getCompressor() {
if (codecPoolSupportsConf) {
return CodecPool.getCompressor(this, getConf());
} else {
// this is potentially wrong since user's configuration changes between
// different two instances of LzopCodec are not honored.
return CodecPool.getCompressor(this);
}
}
@Override
public String getDefaultExtension() {
return DEFAULT_LZO_EXTENSION;
}
}